ym310_mqtt_api.c 8.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299
  1. /*
  2. * YM310 MQTT API Layer - SIMPLE FIX v3.3.1
  3. *
  4. * Fixes in this version:
  5. * - Simple fix: publish failure triggers immediate reconnect
  6. *
  7. * Version: 3.3.1
  8. */
  9. #define DBG_TAG "Y-API"
  10. #define DBG_LVL DBG_INFO
  11. #include <rtdbg.h>
  12. #include "ym310_mqtt.h"
  13. #include <string.h>
  14. #include <stdio.h>
  15. /* External declaration for the new async with data function */
  16. extern at_resp_type_t ym310_at_exec_async_with_data(const char *cmd, const char *data, rt_uint16_t data_len,
  17. const char *urc_key, char *urc_buf, rt_uint16_t urc_size,
  18. rt_int32_t timeout_ms);
  19. /*=============================================================================
  20. * FIXED: Internal Send Function - simple fix
  21. *============================================================================*/
  22. static int ym310_mqtt_send_internal(rt_uint8_t client_idx, const char *topic,
  23. const char *payload, rt_uint16_t len,
  24. rt_uint8_t qos)
  25. {
  26. at_resp_type_t resp;
  27. char cmd[512];
  28. char urc_buf[64];
  29. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  30. rt_uint8_t msg_id = 1;
  31. rt_uint16_t payload_len;
  32. int retry_count = 0;
  33. #define MAX_SEND_RETRIES 2
  34. if (!inst->config.enabled)
  35. {
  36. LOG_E("MQTT[%d] Send: Instance not enabled", client_idx);
  37. return -RT_ERROR;
  38. }
  39. if (inst->state != MQTT_STATE_CONNECTED)
  40. {
  41. LOG_E("MQTT[%d] Send: Not connected (state=%d)", client_idx, inst->state);
  42. return -RT_ERROR;
  43. }
  44. /* Determine payload length */
  45. if (len == 0)
  46. {
  47. payload_len = rt_strlen(payload);
  48. }
  49. else
  50. {
  51. payload_len = len;
  52. }
  53. /* Check payload size */
  54. if (payload_len > 1600)
  55. {
  56. LOG_E("MQTT[%d] Send: Payload too large (%d > 1600)",
  57. client_idx, payload_len);
  58. return -RT_ERROR;
  59. }
  60. LOG_D("MQTT[%d] Sending %d bytes to %s", client_idx, payload_len, topic);
  61. /* FIXED: Proper result parsing and retry logic */
  62. do
  63. {
  64. /* Format: AT+QMTPUBEX=client_idx,msgid,qos,retain,"topic",length */
  65. rt_snprintf(cmd, sizeof(cmd), "AT+QMTPUBEX=%d,%d,%d,0,\"%s\",%d",
  66. client_idx, msg_id, qos, topic, payload_len);
  67. /* Use the async_with_data function */
  68. resp = ym310_at_exec_async_with_data(cmd, payload, payload_len,
  69. "+QMTPUBEX:", urc_buf, sizeof(urc_buf), 30000);
  70. if (resp == AT_RESP_OK)
  71. {
  72. /* FIXED: Parse result correctly - +QMTPUBEX: client_idx,msgid,result
  73. * client_idx is at index 0, msgid at index 1, result at index 2 */
  74. int result = ym310_str_extract_nth_int(urc_buf, 2);
  75. if (result == 0)
  76. {
  77. /* Success! Exit retry loop */
  78. break;
  79. }
  80. else if (result > 0)
  81. {
  82. /* Some modules return non-zero even on success, check if message was actually sent */
  83. /* Result values: 0=success, 1=maybe success, >1=error */
  84. if (result == 1)
  85. {
  86. /* Result=1 might mean success for some modules - check if AT response was OK */
  87. LOG_D("MQTT[%d] Publish result=%d (possible success)", client_idx, result);
  88. break; /* Treat result=1 as success */
  89. }
  90. LOG_W("MQTT[%d] Publish result=%d, retrying (%d/%d)...",
  91. client_idx, result, retry_count + 1, MAX_SEND_RETRIES);
  92. }
  93. else
  94. {
  95. /* Failed to parse result */
  96. LOG_W("MQTT[%d] Publish result parse error, retrying (%d/%d)...",
  97. client_idx, retry_count + 1, MAX_SEND_RETRIES);
  98. }
  99. }
  100. else
  101. {
  102. LOG_W("MQTT[%d] Publish resp=%d, retrying (%d/%d)...",
  103. client_idx, resp, retry_count + 1, MAX_SEND_RETRIES);
  104. }
  105. retry_count++;
  106. if (retry_count < MAX_SEND_RETRIES)
  107. {
  108. rt_thread_mdelay(100 * retry_count);
  109. }
  110. } while (retry_count < MAX_SEND_RETRIES);
  111. /* Check final result */
  112. if (resp != AT_RESP_OK)
  113. {
  114. LOG_E("MQTT[%d] Publish failed after %d retries: resp=%d",
  115. client_idx, MAX_SEND_RETRIES, resp);
  116. /* SIMPLE FIX: Just update error count and trigger reconnect */
  117. inst->err_count++;
  118. /* 发布失败,立即触发重连 */
  119. LOG_W("MQTT[%d] Publish failed, triggering reconnect", client_idx);
  120. ym310_mqtt_force_reconnect_all();
  121. return -RT_ERROR;
  122. }
  123. /* Update statistics */
  124. inst->tx_count++;
  125. inst->last_activity = rt_tick_get();
  126. LOG_D("MQTT[%d] Message sent successfully", client_idx);
  127. g_last_mqtt_success_tick = rt_tick_get();
  128. return RT_EOK;
  129. }
  130. /*=============================================================================
  131. * Public Send API - unchanged
  132. *============================================================================*/
  133. int ym310_mqtt_send_primary(const char *payload, rt_uint16_t len)
  134. {
  135. int ret;
  136. /* 1. 检查 MQTT 实例是否已连接(快速状态检查) */
  137. if (!ym310_mqtt_is_connected(0))
  138. {
  139. LOG_W("MQTT[0] Not connected, cannot send");
  140. return -RT_ERROR;
  141. }
  142. /* 2. 调用内部发送函数(已包含重试和错误处理) */
  143. ret = ym310_mqtt_send_internal(0, g_mqtt_instances[0].config.pub_topic,
  144. payload, len, g_mqtt_instances[0].config.qos);
  145. return ret;
  146. }
  147. int ym310_mqtt_send_secondary(const char *payload, rt_uint16_t len)
  148. {
  149. int ret;
  150. /* 1. 检查 MQTT 实例是否已连接(快速状态检查) */
  151. if (!ym310_mqtt_is_connected(1))
  152. {
  153. LOG_W("MQTT[0] Not connected, cannot send");
  154. return -RT_ERROR;
  155. }
  156. /* 2. 检查网络底层是否连通(可能阻塞,但发送本身会持有AT互斥锁,可接受) */
  157. if (!ym310_check_network_connected())
  158. {
  159. LOG_W("Network disconnected, triggering reconnect and abort send");
  160. /* 触发 MQTT[0] 重连(异步) */
  161. ym310_mqtt_force_reconnect(1);
  162. return -RT_ERROR;
  163. }
  164. /* 3. 调用内部发送函数(已包含重试和错误处理) */
  165. ret = ym310_mqtt_send_internal(1, g_mqtt_instances[1].config.pub_topic,
  166. payload, len, g_mqtt_instances[1].config.qos);
  167. return ret;
  168. }
  169. int ym310_mqtt_send_to(rt_uint8_t client_idx, const char *topic,
  170. const char *payload, rt_uint16_t len, rt_uint8_t qos)
  171. {
  172. int ret;
  173. if (client_idx >= MQTT_INSTANCE_MAX)
  174. {
  175. LOG_E("Invalid client_idx: %d", client_idx);
  176. return -RT_ERROR;
  177. }
  178. if (topic == RT_NULL || payload == RT_NULL)
  179. {
  180. LOG_E("Invalid parameters");
  181. return -RT_ERROR;
  182. }
  183. /* Check connection state - no mutex needed for simple state check */
  184. if (!ym310_mqtt_is_connected(client_idx))
  185. {
  186. LOG_W("MQTT[%d] Not connected, cannot send", client_idx);
  187. return -RT_ERROR;
  188. }
  189. /* AT command functions handle their own mutex internally */
  190. ret = ym310_mqtt_send_internal(client_idx, topic, payload, len, qos);
  191. return ret;
  192. }
  193. /*=============================================================================
  194. * Statistics API (unchanged)
  195. *============================================================================*/
  196. int ym310_mqtt_get_stats(rt_uint8_t client_idx, rt_uint32_t *tx_count,
  197. rt_uint32_t *rx_count, rt_uint32_t *err_count)
  198. {
  199. if (client_idx >= MQTT_INSTANCE_MAX)
  200. {
  201. return -RT_ERROR;
  202. }
  203. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  204. if (tx_count != RT_NULL)
  205. *tx_count = inst->tx_count;
  206. if (rx_count != RT_NULL)
  207. *rx_count = inst->rx_count;
  208. if (err_count != RT_NULL)
  209. *err_count = inst->err_count;
  210. return RT_EOK;
  211. }
  212. int ym310_mqtt_reset_stats(rt_uint8_t client_idx)
  213. {
  214. if (client_idx >= MQTT_INSTANCE_MAX)
  215. {
  216. return -RT_ERROR;
  217. }
  218. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  219. inst->tx_count = 0;
  220. inst->rx_count = 0;
  221. inst->err_count = 0;
  222. LOG_I("MQTT[%d] Statistics reset", client_idx);
  223. return RT_EOK;
  224. }
  225. int ym310_mqtt_update_config(rt_uint8_t client_idx, mqtt_config_t *config)
  226. {
  227. if (client_idx >= MQTT_INSTANCE_MAX || config == RT_NULL)
  228. {
  229. return -RT_ERROR;
  230. }
  231. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  232. rt_memcpy(&inst->config, config, sizeof(mqtt_config_t));
  233. LOG_I("MQTT[%d] Configuration updated", client_idx);
  234. return RT_EOK;
  235. }
  236. int ym310_mqtt_get_config(rt_uint8_t client_idx, mqtt_config_t *config)
  237. {
  238. if (client_idx >= MQTT_INSTANCE_MAX || config == RT_NULL)
  239. {
  240. return -RT_ERROR;
  241. }
  242. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  243. rt_memcpy(config, &inst->config, sizeof(mqtt_config_t));
  244. return RT_EOK;
  245. }