| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299 |
- /*
- * YM310 MQTT API Layer - SIMPLE FIX v3.3.1
- *
- * Fixes in this version:
- * - Simple fix: publish failure triggers immediate reconnect
- *
- * Version: 3.3.1
- */
- #define DBG_TAG "Y-API"
- #define DBG_LVL DBG_INFO
- #include <rtdbg.h>
- #include "ym310_mqtt.h"
- #include <string.h>
- #include <stdio.h>
- /* External declaration for the new async with data function */
- extern at_resp_type_t ym310_at_exec_async_with_data(const char *cmd, const char *data, rt_uint16_t data_len,
- const char *urc_key, char *urc_buf, rt_uint16_t urc_size,
- rt_int32_t timeout_ms);
- /*=============================================================================
- * FIXED: Internal Send Function - simple fix
- *============================================================================*/
- static int ym310_mqtt_send_internal(rt_uint8_t client_idx, const char *topic,
- const char *payload, rt_uint16_t len,
- rt_uint8_t qos)
- {
- at_resp_type_t resp;
- char cmd[512];
- char urc_buf[64];
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- rt_uint8_t msg_id = 1;
- rt_uint16_t payload_len;
- int retry_count = 0;
- #define MAX_SEND_RETRIES 2
- if (!inst->config.enabled)
- {
- LOG_E("MQTT[%d] Send: Instance not enabled", client_idx);
- return -RT_ERROR;
- }
- if (inst->state != MQTT_STATE_CONNECTED)
- {
- LOG_E("MQTT[%d] Send: Not connected (state=%d)", client_idx, inst->state);
- return -RT_ERROR;
- }
- /* Determine payload length */
- if (len == 0)
- {
- payload_len = rt_strlen(payload);
- }
- else
- {
- payload_len = len;
- }
- /* Check payload size */
- if (payload_len > 1600)
- {
- LOG_E("MQTT[%d] Send: Payload too large (%d > 1600)",
- client_idx, payload_len);
- return -RT_ERROR;
- }
- LOG_D("MQTT[%d] Sending %d bytes to %s", client_idx, payload_len, topic);
- /* FIXED: Proper result parsing and retry logic */
- do
- {
- /* Format: AT+QMTPUBEX=client_idx,msgid,qos,retain,"topic",length */
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTPUBEX=%d,%d,%d,0,\"%s\",%d",
- client_idx, msg_id, qos, topic, payload_len);
- /* Use the async_with_data function */
- resp = ym310_at_exec_async_with_data(cmd, payload, payload_len,
- "+QMTPUBEX:", urc_buf, sizeof(urc_buf), 30000);
- if (resp == AT_RESP_OK)
- {
- /* FIXED: Parse result correctly - +QMTPUBEX: client_idx,msgid,result
- * client_idx is at index 0, msgid at index 1, result at index 2 */
- int result = ym310_str_extract_nth_int(urc_buf, 2);
- if (result == 0)
- {
- /* Success! Exit retry loop */
- break;
- }
- else if (result > 0)
- {
- /* Some modules return non-zero even on success, check if message was actually sent */
- /* Result values: 0=success, 1=maybe success, >1=error */
- if (result == 1)
- {
- /* Result=1 might mean success for some modules - check if AT response was OK */
- LOG_D("MQTT[%d] Publish result=%d (possible success)", client_idx, result);
- break; /* Treat result=1 as success */
- }
- LOG_W("MQTT[%d] Publish result=%d, retrying (%d/%d)...",
- client_idx, result, retry_count + 1, MAX_SEND_RETRIES);
- }
- else
- {
- /* Failed to parse result */
- LOG_W("MQTT[%d] Publish result parse error, retrying (%d/%d)...",
- client_idx, retry_count + 1, MAX_SEND_RETRIES);
- }
- }
- else
- {
- LOG_W("MQTT[%d] Publish resp=%d, retrying (%d/%d)...",
- client_idx, resp, retry_count + 1, MAX_SEND_RETRIES);
- }
- retry_count++;
- if (retry_count < MAX_SEND_RETRIES)
- {
- rt_thread_mdelay(100 * retry_count);
- }
- } while (retry_count < MAX_SEND_RETRIES);
- /* Check final result */
- if (resp != AT_RESP_OK)
- {
- LOG_E("MQTT[%d] Publish failed after %d retries: resp=%d",
- client_idx, MAX_SEND_RETRIES, resp);
- /* SIMPLE FIX: Just update error count and trigger reconnect */
- inst->err_count++;
- /* 发布失败,立即触发重连 */
- LOG_W("MQTT[%d] Publish failed, triggering reconnect", client_idx);
- ym310_mqtt_force_reconnect_all();
- return -RT_ERROR;
- }
- /* Update statistics */
- inst->tx_count++;
- inst->last_activity = rt_tick_get();
- LOG_D("MQTT[%d] Message sent successfully", client_idx);
- g_last_mqtt_success_tick = rt_tick_get();
- return RT_EOK;
- }
- /*=============================================================================
- * Public Send API - unchanged
- *============================================================================*/
- int ym310_mqtt_send_primary(const char *payload, rt_uint16_t len)
- {
- int ret;
- /* 1. 检查 MQTT 实例是否已连接(快速状态检查) */
- if (!ym310_mqtt_is_connected(0))
- {
- LOG_W("MQTT[0] Not connected, cannot send");
- return -RT_ERROR;
- }
- /* 2. 调用内部发送函数(已包含重试和错误处理) */
- ret = ym310_mqtt_send_internal(0, g_mqtt_instances[0].config.pub_topic,
- payload, len, g_mqtt_instances[0].config.qos);
- return ret;
- }
- int ym310_mqtt_send_secondary(const char *payload, rt_uint16_t len)
- {
- int ret;
- /* 1. 检查 MQTT 实例是否已连接(快速状态检查) */
- if (!ym310_mqtt_is_connected(1))
- {
- LOG_W("MQTT[0] Not connected, cannot send");
- return -RT_ERROR;
- }
- /* 2. 检查网络底层是否连通(可能阻塞,但发送本身会持有AT互斥锁,可接受) */
- if (!ym310_check_network_connected())
- {
- LOG_W("Network disconnected, triggering reconnect and abort send");
- /* 触发 MQTT[0] 重连(异步) */
- ym310_mqtt_force_reconnect(1);
- return -RT_ERROR;
- }
- /* 3. 调用内部发送函数(已包含重试和错误处理) */
- ret = ym310_mqtt_send_internal(1, g_mqtt_instances[1].config.pub_topic,
- payload, len, g_mqtt_instances[1].config.qos);
- return ret;
- }
- int ym310_mqtt_send_to(rt_uint8_t client_idx, const char *topic,
- const char *payload, rt_uint16_t len, rt_uint8_t qos)
- {
- int ret;
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- LOG_E("Invalid client_idx: %d", client_idx);
- return -RT_ERROR;
- }
- if (topic == RT_NULL || payload == RT_NULL)
- {
- LOG_E("Invalid parameters");
- return -RT_ERROR;
- }
- /* Check connection state - no mutex needed for simple state check */
- if (!ym310_mqtt_is_connected(client_idx))
- {
- LOG_W("MQTT[%d] Not connected, cannot send", client_idx);
- return -RT_ERROR;
- }
- /* AT command functions handle their own mutex internally */
- ret = ym310_mqtt_send_internal(client_idx, topic, payload, len, qos);
- return ret;
- }
- /*=============================================================================
- * Statistics API (unchanged)
- *============================================================================*/
- int ym310_mqtt_get_stats(rt_uint8_t client_idx, rt_uint32_t *tx_count,
- rt_uint32_t *rx_count, rt_uint32_t *err_count)
- {
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- return -RT_ERROR;
- }
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- if (tx_count != RT_NULL)
- *tx_count = inst->tx_count;
- if (rx_count != RT_NULL)
- *rx_count = inst->rx_count;
- if (err_count != RT_NULL)
- *err_count = inst->err_count;
- return RT_EOK;
- }
- int ym310_mqtt_reset_stats(rt_uint8_t client_idx)
- {
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- return -RT_ERROR;
- }
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- inst->tx_count = 0;
- inst->rx_count = 0;
- inst->err_count = 0;
- LOG_I("MQTT[%d] Statistics reset", client_idx);
- return RT_EOK;
- }
- int ym310_mqtt_update_config(rt_uint8_t client_idx, mqtt_config_t *config)
- {
- if (client_idx >= MQTT_INSTANCE_MAX || config == RT_NULL)
- {
- return -RT_ERROR;
- }
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- rt_memcpy(&inst->config, config, sizeof(mqtt_config_t));
- LOG_I("MQTT[%d] Configuration updated", client_idx);
- return RT_EOK;
- }
- int ym310_mqtt_get_config(rt_uint8_t client_idx, mqtt_config_t *config)
- {
- if (client_idx >= MQTT_INSTANCE_MAX || config == RT_NULL)
- {
- return -RT_ERROR;
- }
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- rt_memcpy(config, &inst->config, sizeof(mqtt_config_t));
- return RT_EOK;
- }
|