/* * 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 #include "ym310_mqtt.h" #include #include /* 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; }