/* mqtt_client.c - 精简版 MQTT 3.1.1 客户端(ESP-01S TCP 透传之上) * * 只做设备需要的最小集:CONNECT / SUBSCRIBE(QoS0) / PUBLISH(QoS0) / PINGREQ, * 以及 CONNACK / SUBACK / PINGRESP / PUBLISH 的解析。 * 服务器下发的 PUBLISH payload(JSON) 直接交给 network.c 的 net_dispatch_message()。 */ #include "mqtt_client.h" #include "esp8266.h" #include "air780e.h" /* U2_CopyBuff / bg_delay */ #include "network.h" #include "log.h" #include "main.h" #include "string.h" #define MQTT_TX_BUF_SIZE 512 /* 单发送缓冲:CONNECT/SUB/PUBLISH 最大约 350B */ #define MQTT_RX_JSON_SIZE 384 /* 下行 PUBLISH payload 暂存 */ static uint8_t tx_buf[MQTT_TX_BUF_SIZE]; static char s_cid[24]; static char s_user[24]; static char s_pass[24]; static volatile uint8_t s_connack_ok = 0; static volatile uint8_t s_suback_ok = 0; volatile uint8_t g_mqttc_ping_outstanding = 0; void mqttc_init(const char *client_id, const char *user, const char *pass) { strncpy(s_cid, client_id, sizeof(s_cid) - 1); s_cid[sizeof(s_cid)-1] = 0; strncpy(s_user, user, sizeof(s_user) - 1); s_user[sizeof(s_user)-1] = 0; strncpy(s_pass, pass, sizeof(s_pass) - 1); s_pass[sizeof(s_pass)-1] = 0; s_connack_ok = 0; s_suback_ok = 0; g_mqttc_ping_outstanding = 0; } /*==== 剩余长度编码(MQTT 变长编码,我们的报文 < 16KB,最多 3 字节)====*/ static uint8_t encode_rl(uint8_t *out, uint32_t len) { uint8_t n = 0; do { uint8_t b = len % 128; len /= 128; if (len > 0) b |= 0x80; out[n++] = b; } while (len > 0); return n; } /* 追加一个带 2 字节长度前缀的 UTF-8 字符串,返回新偏移 */ static uint16_t put_str(uint8_t *buf, uint16_t pos, const char *s) { uint16_t l = strlen(s); buf[pos++] = (uint8_t)(l >> 8); buf[pos++] = (uint8_t)(l & 0xFF); memcpy(buf + pos, s, l); return pos + l; } /*==== CONNECT:clean session + username + password ====*/ int mqttc_connect(void) { uint16_t pos, var_pos, rl_pos; uint8_t rl_len; tx_buf[0] = 0x10; var_pos = 1 + 3; /* 预留剩余长度(最多3字节),最后回填 */ /* 可变头: "MQTT" + 协议级别4 + 连接标志0xC2 + keepalive */ var_pos = put_str(tx_buf, var_pos, "MQTT"); tx_buf[var_pos++] = 0x04; tx_buf[var_pos++] = 0xC2; tx_buf[var_pos++] = (uint8_t)(MQTT_KEEPALIVE_S >> 8); tx_buf[var_pos++] = (uint8_t)(MQTT_KEEPALIVE_S & 0xFF); /* 载荷: ClientID / Username / Password */ var_pos = put_str(tx_buf, var_pos, s_cid); var_pos = put_str(tx_buf, var_pos, s_user); var_pos = put_str(tx_buf, var_pos, s_pass); /* 回填剩余长度(把可变头整体后移对齐到编码长度) */ pos = var_pos - (1 + 3); rl_len = encode_rl(tx_buf + 1, pos); if (rl_len < 3) { memmove(tx_buf + 1 + rl_len, tx_buf + 1 + 3, pos); var_pos -= (3 - rl_len); } rl_pos = var_pos; log_info("> WiFi MQTT: CONNECT cid=%s", s_cid); esp_send(tx_buf, rl_pos); /* 等 CONNACK,最多 5s,期间切片喂狗 */ s_connack_ok = 0; for (int i = 0; i < 50; i++) { bg_delay(100); while (mqttc_process() == 1) { } if (s_connack_ok) { log_info("> WiFi MQTT: CONNACK 成功"); return 0; } } log_warn("> WiFi MQTT: CONNACK 超时"); return -1; } /*==== SUBSCRIBE QoS0,包标识固定 1 ====*/ int mqttc_subscribe(const char *topic) { uint16_t pos = 0; tx_buf[pos++] = 0x82; /* 剩余长度 = 2(包标识) + 2 + topic + 1(QoS) */ pos += encode_rl(tx_buf + pos, (uint32_t)(2 + 2 + strlen(topic) + 1)); tx_buf[pos++] = 0x00; tx_buf[pos++] = 0x01; pos = put_str(tx_buf, pos, topic); tx_buf[pos++] = 0x00; log_info("> WiFi MQTT: SUB %s", topic); esp_send(tx_buf, pos); s_suback_ok = 0; for (int i = 0; i < 50; i++) { bg_delay(100); while (mqttc_process() == 1) { } if (s_suback_ok) { log_info("> WiFi MQTT: SUBACK 成功"); return 0; } } log_warn("> WiFi MQTT: SUBACK 超时"); return -1; } /*==== PUBLISH QoS0 ====*/ int mqttc_publish(const char *topic, const char *payload) { uint16_t pos = 0; uint32_t rl = 2 + strlen(topic) + strlen(payload); if (rl + 5 > MQTT_TX_BUF_SIZE) { log_warn("> WiFi MQTT: publish 报文过长 (%u)", (unsigned)rl); return -1; } tx_buf[pos++] = 0x30; pos += encode_rl(tx_buf + pos, rl); pos = put_str(tx_buf, pos, topic); memcpy(tx_buf + pos, payload, strlen(payload)); pos += strlen(payload); esp_send(tx_buf, pos); return 0; } /*==== QoS1 发布 + PUBACK 等待(仅平台指令应答用, 确认 broker 收到) ====*/ static volatile uint8_t s_puback_ok = 0; static volatile uint16_t s_puback_id = 0; int mqttc_publish_qos1(const char *topic, const char *payload, uint16_t pkt_id) { uint16_t pos = 0; uint32_t rl = 2 + strlen(topic) + 2 + strlen(payload); /* +2 为报文标识符 */ if (rl + 5 > MQTT_TX_BUF_SIZE) { log_warn("> WiFi MQTT: publish 报文过长 (%u)", (unsigned)rl); return -1; } tx_buf[pos++] = 0x32; /* PUBLISH, QoS1 */ pos += encode_rl(tx_buf + pos, rl); pos = put_str(tx_buf, pos, topic); tx_buf[pos++] = (uint8_t)(pkt_id >> 8); tx_buf[pos++] = (uint8_t)(pkt_id & 0xFF); memcpy(tx_buf + pos, payload, strlen(payload)); pos += strlen(payload); s_puback_id = pkt_id; s_puback_ok = 0; esp_send(tx_buf, pos); return 0; } /* 等 PUBACK: 收包解析在 net 任务(mqttc_process), 这里只轮询标志 */ int mqttc_wait_puback(uint16_t pkt_id, int timeout_ms) { while (timeout_ms > 0) { if (s_puback_ok && s_puback_id == pkt_id) return 0; bg_delay(20); timeout_ms -= 20; } return -1; } void mqttc_ping(void) { tx_buf[0] = 0xC0; tx_buf[1] = 0x00; esp_send(tx_buf, 2); g_mqttc_ping_outstanding = 1; } /*==== 接收解析 ====*/ /* 二进制安全的子串查找(透传数据里可能含 0x00,如 PINGRESP 0xD0 0x00) */ static uint8_t *memfind(uint8_t *hay, uint16_t hay_len, const char *needle) { uint16_t nlen = strlen(needle); if (hay_len < nlen) return NULL; for (uint16_t i = 0; i <= hay_len - nlen; i++) { if (memcmp(hay + i, needle, nlen) == 0) return hay + i; } return NULL; } /* 消费缓冲前 n 字节 */ static void consume(uint16_t n) { if (n >= U2_CopyIndex) { U2_CopyIndex = 0; } else { memmove(U2_CopyBuff, U2_CopyBuff + n, U2_CopyIndex - n); U2_CopyIndex -= n; } U2_CopyBuff[U2_CopyIndex] = '\0'; } int mqttc_process(void) { uint8_t *buf = U2_CopyBuff; uint16_t len = U2_CopyIndex; if (len == 0) { U2_CopyFlag = 0; return 0; } /* 断链检测:透传模式下 TCP 断开 ESP 会输出 "CLOSED" */ if (memfind(buf, len, "CLOSED")) { consume(U2_CopyIndex); U2_CopyFlag = 0; return -1; } /* 帧同步:首字节高 4 位必须是合法报文类型(1~14),否则逐字节丢弃 */ uint8_t type = buf[0] >> 4; if (type < 1 || type > 14) { consume(1); return 1; } /* 解码剩余长度(最多 4 字节) */ uint32_t rl = 0, mult = 1; uint8_t rl_bytes = 0; for (uint8_t i = 1; i <= 4; i++) { if (i >= len) return 0; /* 头还没收全,等数据 */ rl += (uint32_t)(buf[i] & 0x7F) * mult; rl_bytes = i; if (!(buf[i] & 0x80)) break; mult *= 128; } if (rl_bytes == 0) return 0; uint32_t total = 1 + rl_bytes + rl; if (total > U2_COPY_SIZE - 1) { /* 异常巨型报文,丢弃首字节重新同步 */ consume(1); return 1; } if (len < total) return 0; /* 报文没收全,等数据 */ uint8_t *pkt = buf + 1 + rl_bytes; switch (type) { case 2: /* CONNACK: 0x20 0x02 0x00 rc */ if (rl >= 2 && pkt[1] == 0) s_connack_ok = 1; else log_warn("> WiFi MQTT: CONNACK 被拒绝 rc=%u", rl >= 2 ? pkt[1] : 0xFF); break; case 9: /* SUBACK: 0x90 0x03 pktIdHi pktIdLo rc,rc=0x80 表示订阅被拒绝 */ if (rl >= 3 && pkt[2] != 0x80) s_suback_ok = 1; else log_warn("> WiFi MQTT: SUBACK 被拒绝"); break; case 4: /* PUBACK: 0x40 0x02 pktIdHi pktIdLo */ if (rl >= 2) { uint16_t id = (uint16_t)(((uint16_t)pkt[0] << 8) | pkt[1]); if (id == s_puback_id) s_puback_ok = 1; } break; case 13: /* PINGRESP */ g_mqttc_ping_outstanding = 0; break; case 3: { /* PUBLISH QoS0:取 payload(JSON) 交网络层分发 */ static char json[MQTT_RX_JSON_SIZE]; uint16_t tlen = ((uint16_t)pkt[0] << 8) | pkt[1]; if (2 + tlen <= rl) { uint32_t plen = rl - 2 - tlen; if (plen >= sizeof(json)) plen = sizeof(json) - 1; memcpy(json, pkt + 2 + tlen, plen); json[plen] = '\0'; log_info("> WiFi MQTT 接收: %.200s", json); net_dispatch_message(json); } break; } default: /* PUBACK 等其他报文,忽略 */ break; } consume((uint16_t)total); return 1; }