root cause 1: 3次失败断连后, iot_connect_broker() 只设
g_iot_socket=SocketId_TCP, 但 WCHNET_SocketClose 后
socket 已销毁, 没有重建 → TCP connect 一直超时
root cause 2: old socket close 未完成就立即重连,
WCHNET_SocketCreat 可能失败
fix:
1. iot_connect_broker() 先调 WCHNET_CreateTcpMqttSocket()
真正 create + connect
2. 3次失败后加 3s 延迟再重连, 给 WCHNET 清理时间
990 lines
42 KiB
C
990 lines
42 KiB
C
/**
|
||
******************************************************************************
|
||
* @file iot_mqtt_srv.c
|
||
* @author wangfq
|
||
* @version V1.0
|
||
* @date 2026-07-03
|
||
* @brief IoT MQTT 客户端实现 — DLD960 IoT 协议
|
||
*
|
||
* 连接流程: TCP connect → MQTT CONNECT → CONNACK → SUBSCRIBE → SUBACK → READY
|
||
* 断线重连: 指数退避 5s~60s
|
||
* 数据上报: loop_data / event_report / heartbeat
|
||
******************************************************************************
|
||
*/
|
||
|
||
#include "CONFIG.h"
|
||
#include "iot_mqtt_srv.h"
|
||
#include "eth_driver.h"
|
||
#include "wchnet.h"
|
||
#include "net_srv.h"
|
||
#include "MQTTPacket.h"
|
||
#include "cmcng.h"
|
||
#include "loop_uart_proto.h"
|
||
#include "simple_json.h"
|
||
#include "tcp_json_srv.h"
|
||
#include "storage.h"
|
||
#include "ch32v20x_iwdg.h"
|
||
#include <string.h>
|
||
#include <stdlib.h>
|
||
#include <stdio.h>
|
||
|
||
/*===========================================================================
|
||
* Global State
|
||
*===========================================================================*/
|
||
|
||
|
||
extern ReportConfig g_report_cfg; // from tcp_json_srv.c
|
||
extern uint8_t SocketId_TCP; // from net_srv.c, MQTT socket created by WCHNET_CreateTcpMqttSocket
|
||
uint8_t g_iot_socket = 0xFF; // MQTT TCP socket ID
|
||
IotMqttState g_iot_state = IOT_STATE_DISCONNECTED;
|
||
uint8_t g_iot_msg_id = 0; // MQTT packet identifier
|
||
static uint32_t _iot_last_heartbeat = 0; // 上次 PINGREQ 时刻 (ms), MQTT 保活
|
||
static uint32_t _iot_reconnect_deadline = 0;
|
||
static uint32_t _iot_reconnect_backoff = 0;
|
||
static uint32_t _iot_connect_start = 0; // TCP connect 开始时刻
|
||
|
||
/*===========================================================================
|
||
* 传感器数据缓存 & 上报间隔控制
|
||
*===========================================================================*/
|
||
#define IOT_MQTT_FAST_INTERVAL_MS 300 // 变化态快速上报间隔 (ms)
|
||
#define IOT_MQTT_VARIATION_THRESHOLD 10 // diff 阈值: >= 此值进入快速上报
|
||
#define IOT_MQTT_IDLE_INTERVAL_MIN_MS 1000 // 空闲间隔最小 1s (防 interval=0 死锁)
|
||
|
||
static LUP_SensorReport _cached_sr; // 最新传感器数据缓存
|
||
static uint8_t _cached_sr_valid; // 缓存有效标志
|
||
static uint32_t _last_publish_ms; // 上次上报时刻 (ms)
|
||
|
||
/*===========================================================================
|
||
* 事件上报 event_report — 协议 V1.04: 平台必答 + 设备重发
|
||
*
|
||
* 检测: car_state/loop_state 翻转沿 (每帧比对, 经 lup 回调 + Step1 双路汇聚)
|
||
* 队列: 16 深环形队列, 溢出丢最旧; 事件仅在收到 ACK 后出队
|
||
* 重发: 5s 超时, 同 msg_id/原始 ts, 最多 3 次; 耗尽后挂起,
|
||
* 待 MQTT 重连或新事件到达时以新 msg_id 合并重报
|
||
* 应答: 平台回 {cmd:"event_report", msg_id(回显), code}, 由
|
||
* manage_mqtt_recv_message (net_srv.c) 路由到 iot_evt_handle_ack()
|
||
*===========================================================================*/
|
||
#define IOT_EVT_QUEUE_DEPTH 16 // 待发队列深度 (协议建议 >=16)
|
||
#define IOT_EVT_MAX_PER_PKT 6 // 单包最多事件数 (6条约358B < 500B发布上限)
|
||
#define IOT_EVT_ACK_TIMEOUT_MS 5000 // 确认超时 (协议 V1.04)
|
||
#define IOT_EVT_MAX_RETRY 3 // 最大重发次数 (含首发共4次)
|
||
|
||
enum { IOT_EVT_CAR_ENTER = 0, IOT_EVT_CAR_LEAVE, IOT_EVT_LOOP_CUT, IOT_EVT_LOOP_RESTORE };
|
||
|
||
typedef struct {
|
||
uint8_t type; // IOT_EVT_xxx
|
||
uint8_t ch; // 通道号 1-based
|
||
uint32_t value; // car_leave=通过时间(50ms), loop_restore=断开时长(50ms)
|
||
} IotEvent;
|
||
|
||
static IotEvent _evt_queue[IOT_EVT_QUEUE_DEPTH]; // 环形队列
|
||
static uint8_t _evt_head, _evt_count;
|
||
|
||
static uint8_t _evt_prev_car[LUP_COIL_COUNT]; // 上一帧 car_state
|
||
static uint8_t _evt_prev_loop[LUP_COIL_COUNT]; // 上一帧 loop_state (1=断开)
|
||
static uint32_t _evt_cut_ms[LUP_COIL_COUNT]; // 断开起始时刻 (mstick)
|
||
static uint8_t _evt_prev_valid; // 首帧只建快照不出事件
|
||
|
||
static uint32_t _evt_msg_id; // event_report 独立 msg_id (uint32 递增)
|
||
static uint32_t _evt_pend_id; // 待确认包 msg_id, 0=无
|
||
static uint32_t _evt_pend_ts; // 待确认包 ts (首发时刻, 重发不刷新)
|
||
static uint8_t _evt_pend_n; // 待确认包内事件数 (= 队列头 N 条)
|
||
static uint32_t _evt_sent_ms; // 本次发送时刻 (超时计时)
|
||
static uint8_t _evt_retry; // 已重发次数
|
||
static uint8_t _evt_gaveup; // 重试耗尽挂起标志 (新事件/重连时解除)
|
||
|
||
/* 入队: 满则丢最旧; 若丢掉的事件属于待确认包, 作废该包 (事件仍按新包重报) */
|
||
static void iot_evt_enqueue(uint8_t type, uint8_t ch, uint32_t value) {
|
||
if (_evt_count >= IOT_EVT_QUEUE_DEPTH) {
|
||
_evt_head = (_evt_head + 1) % IOT_EVT_QUEUE_DEPTH;
|
||
_evt_count--;
|
||
if (_evt_pend_id) { // 待确认包引用队列头, 头被丢弃则作废重来
|
||
_evt_pend_id = 0;
|
||
if (_evt_pend_n) _evt_pend_n = 0;
|
||
}
|
||
PRINT("EVT: queue overflow, drop oldest\n");
|
||
}
|
||
{
|
||
IotEvent *e = &_evt_queue[(_evt_head + _evt_count) % IOT_EVT_QUEUE_DEPTH];
|
||
e->type = type; e->ch = ch; e->value = value;
|
||
_evt_count++;
|
||
}
|
||
_evt_gaveup = 0; // 新事件到达 → 解除挂起, 触发合并上报
|
||
PRINT("EVT: enqueue type=%d ch=%d val=%lu (cnt=%d)\n", type, ch, value, _evt_count);
|
||
}
|
||
|
||
/* 沿检测: 每收到一帧有效 0xC0 调用一次 (两条消费路径都要喂) */
|
||
static void iot_evt_feed(const LUP_SensorReport *sr) {
|
||
uint8_t i;
|
||
if (!_evt_prev_valid) { // 首帧: 只建快照 (上电已有车不算进入沿)
|
||
for (i = 0; i < LUP_COIL_COUNT; i++) {
|
||
_evt_prev_car[i] = (i < sr->coil_count) ? sr->coils[i].car_state : 0;
|
||
_evt_prev_loop[i] = (i < sr->coil_count) ? sr->coils[i].loop_state : 0;
|
||
}
|
||
_evt_prev_valid = 1;
|
||
return;
|
||
}
|
||
for (i = 0; i < sr->coil_count && i < LUP_COIL_COUNT; i++) {
|
||
const LUP_CoilSensor *cs = &sr->coils[i];
|
||
|
||
/* 线圈断开期间 car_state 不可信: 仅前后两帧线圈均正常才判进出车沿 */
|
||
if (!_evt_prev_loop[i] && !cs->loop_state) {
|
||
if (!_evt_prev_car[i] && cs->car_state) {
|
||
iot_evt_enqueue(IOT_EVT_CAR_ENTER, i + 1, 0);
|
||
} else if (_evt_prev_car[i] && !cs->car_state) {
|
||
/* 离开帧 misc_type=时间量 时携带通过时间 (50ms 单位) */
|
||
uint32_t v = (cs->misc_type == 0) ? cs->misc.passtime_ms : 0;
|
||
iot_evt_enqueue(IOT_EVT_CAR_LEAVE, i + 1, v);
|
||
}
|
||
}
|
||
/* 线圈断开/恢复沿 */
|
||
if (!_evt_prev_loop[i] && cs->loop_state) {
|
||
_evt_cut_ms[i] = mstick();
|
||
iot_evt_enqueue(IOT_EVT_LOOP_CUT, i + 1, 0);
|
||
} else if (_evt_prev_loop[i] && !cs->loop_state) {
|
||
iot_evt_enqueue(IOT_EVT_LOOP_RESTORE, i + 1,
|
||
(mstick() - _evt_cut_ms[i]) / 50); // 断开时长, 50ms 单位
|
||
}
|
||
|
||
_evt_prev_car[i] = cs->car_state;
|
||
_evt_prev_loop[i] = cs->loop_state;
|
||
}
|
||
}
|
||
|
||
/* 统一帧摄取: 事件沿检测 + 刷新 loop_data 缓存 (单一数据源)
|
||
两条消费路径 (uart_srv→lup回调 / Step1直读) 都必须走这里。
|
||
教训(2026-07-15): 此前回调只喂事件、缓存只在 Step1 更新, 而 Step1 赢得
|
||
帧竞争的概率实测 <2% → loop_data 发布 34s 前的陈旧快照(车走后 iscar
|
||
仍 true), 且陈旧 diff 锁死 fast_mode 300ms 刷屏。事件与快照必须同源! */
|
||
static void iot_sensor_ingest(const LUP_SensorReport *sr) {
|
||
iot_evt_feed(sr); // 事件沿检测
|
||
memcpy(&_cached_sr, sr, sizeof(LUP_SensorReport)); // 刷新 loop_data 快照
|
||
_cached_sr_valid = 1;
|
||
}
|
||
|
||
/* lup 传感回调: uart_srv 消费路径的帧从这里喂入 (lup_process_frame 已过校验) */
|
||
static void iot_evt_sensor_cb(const uint8_t *pkg, uint16_t len) {
|
||
LUP_SensorReport sr;
|
||
memset(&sr, 0, sizeof(sr));
|
||
if (lup_parse_sensor_report(pkg, len, &sr) == 0) {
|
||
iot_sensor_ingest(&sr);
|
||
}
|
||
}
|
||
|
||
/* 序列化并发布队列头 n 条事件 (首发与重发共用: 同 id/ts 同内容) */
|
||
static void iot_evt_send(uint32_t id, uint32_t ts, uint8_t n) {
|
||
static const char *evt_names[] = {"car_enter", "car_leave", "loop_cut", "loop_restore"};
|
||
static char payload[512]; // static: 避免 2KB 栈溢出
|
||
int w, rem = sizeof(payload);
|
||
char *p = payload;
|
||
uint8_t i;
|
||
|
||
w = snprintf(p, rem, "{\"msg_id\":%lu,\"cmd\":\"event_report\",\"ts\":%lu,"
|
||
"\"data\":{\"events\":[", id, ts);
|
||
if (w < 0 || w >= rem) return;
|
||
p += w; rem -= w;
|
||
|
||
for (i = 0; i < n; i++) {
|
||
const IotEvent *e = &_evt_queue[(_evt_head + i) % IOT_EVT_QUEUE_DEPTH];
|
||
w = snprintf(p, rem, "%s{\"type\":\"%s\",\"ch\":%d,\"value\":%lu}",
|
||
(i > 0) ? "," : "", evt_names[e->type & 0x03], e->ch, e->value);
|
||
if (w < 0 || w >= rem) return;
|
||
p += w; rem -= w;
|
||
}
|
||
snprintf(p, rem, "]}}");
|
||
|
||
mqtt_publish((char *)g_iot_topic.topic_pub, payload, 1);
|
||
PRINT("EVT: publish msg_id=%lu n=%d retry=%d\n", id, n, _evt_retry);
|
||
}
|
||
|
||
/* 平台应答入口 — manage_mqtt_recv_message (net_srv.c) 收到 cmd=event_report 时调用 */
|
||
void iot_evt_handle_ack(uint32_t msg_id, int code) {
|
||
if (_evt_pend_id == 0 || msg_id != _evt_pend_id) {
|
||
PRINT("EVT: ack msg_id=%lu ignored (pend=%lu)\n", msg_id, _evt_pend_id);
|
||
return;
|
||
}
|
||
if (code != 0) { // 平台处理失败 → 视同未确认, 等超时重发
|
||
PRINT("EVT: ack code=%d, keep pending\n", code);
|
||
return;
|
||
}
|
||
/* 确认成功 → 弹出本包事件 */
|
||
_evt_head = (_evt_head + _evt_pend_n) % IOT_EVT_QUEUE_DEPTH;
|
||
_evt_count -= _evt_pend_n;
|
||
PRINT("EVT: ack ok, %d events dequeued (left=%d)\n", _evt_pend_n, _evt_count);
|
||
_evt_pend_id = 0;
|
||
_evt_pend_n = 0;
|
||
_evt_retry = 0;
|
||
}
|
||
|
||
/* 事件发送状态机: 主循环每轮调用 (置于 READY 检查之前, 以感知重连沿) */
|
||
static void iot_evt_process(void) {
|
||
static uint8_t _was_ready = 0;
|
||
uint8_t ready = (g_iot_state == IOT_STATE_READY);
|
||
uint32_t now = mstick();
|
||
|
||
if (ready && !_was_ready) { // 重连沿
|
||
_evt_gaveup = 0;
|
||
if (_evt_pend_id) {
|
||
/* 有未决包 → 用【原 msg_id/原 ts】立即重发 (协议 V1.04 §5.3-2:
|
||
跨重连也须保持同 msg_id, 否则平台 (sn,msg_id) 去重失效致重复入库) */
|
||
_evt_retry = 0; // 重连后重试计数归零, 不吃掉本次
|
||
iot_evt_send(_evt_pend_id, _evt_pend_ts, _evt_pend_n);
|
||
_evt_sent_ms = now;
|
||
_was_ready = ready;
|
||
return;
|
||
}
|
||
}
|
||
_was_ready = ready;
|
||
if (!ready) return;
|
||
|
||
if (_evt_pend_id) { // 有未决包 → 超时重发
|
||
if (now - _evt_sent_ms < IOT_EVT_ACK_TIMEOUT_MS) return;
|
||
if (_evt_retry >= IOT_EVT_MAX_RETRY) {
|
||
PRINT("EVT: retry exhausted msg_id=%lu, hold %d events\n",
|
||
_evt_pend_id, _evt_count);
|
||
_evt_pend_id = 0;
|
||
_evt_gaveup = 1; // 挂起: 新事件或重连时再触发
|
||
return;
|
||
}
|
||
_evt_retry++;
|
||
iot_evt_send(_evt_pend_id, _evt_pend_ts, _evt_pend_n); // 同 id/ts 重发
|
||
_evt_sent_ms = now;
|
||
return;
|
||
}
|
||
|
||
if (_evt_count == 0 || _evt_gaveup) return;
|
||
|
||
/* 发新包: 队列头 N 条合并 */
|
||
_evt_pend_n = (_evt_count > IOT_EVT_MAX_PER_PKT) ? IOT_EVT_MAX_PER_PKT : _evt_count;
|
||
_evt_pend_id = ++_evt_msg_id;
|
||
_evt_pend_ts = dev_time_now(); // 首发时刻 Unix 时间(已校准), 重发保持不刷新
|
||
_evt_retry = 0;
|
||
iot_evt_send(_evt_pend_id, _evt_pend_ts, _evt_pend_n);
|
||
_evt_sent_ms = now;
|
||
}
|
||
|
||
/*===========================================================================
|
||
* Buffers
|
||
*===========================================================================*/
|
||
static uint8_t _iot_recv_buf[IOT_MQTT_RECV_BUF_LEN];
|
||
static uint16_t _iot_recv_len = 0;
|
||
static uint8_t _iot_send_buf[IOT_MQTT_SEND_BUF_LEN];
|
||
static uint8_t _iot_wchnet_buf[RECE_BUF_LEN]; // WCHNET internal recv buffer
|
||
|
||
/*===========================================================================
|
||
* Topic 构建
|
||
*===========================================================================*/
|
||
static char g_iot_dev_serial[13]; // 设备序列码 (12 hex chars)
|
||
|
||
/* 构建 topic: dld960/{sn}/... */
|
||
static int iot_make_topic(char *out, uint16_t out_len, const char *direction,
|
||
const char *category, const char *sub) {
|
||
if (sub) {
|
||
return snprintf(out, out_len, "dld960/%s/%s/%s/%s",
|
||
g_iot_dev_serial, direction, category, sub);
|
||
}
|
||
return snprintf(out, out_len, "dld960/%s/%s/%s",
|
||
g_iot_dev_serial, direction, category);
|
||
}
|
||
|
||
/*===========================================================================
|
||
* MQTT Packet Helpers
|
||
*===========================================================================*/
|
||
|
||
/* 发送 MQTT 报文到 broker — 循环重试防 WCHNET 部分发送
|
||
(2026-07-23: SocketSend 只发 520/607B → 残留拼到下行 → broker 见垃圾 → _raw+RST)
|
||
连续失败追踪: ret=0x11 表示 TCP 发送队列满/连接断开, 3次连续失败 → 主动断连重连
|
||
(TCP超时要等~2分钟, 太慢了, 导致长时间丢数据) */
|
||
static int _iot_send_fail_cnt = 0; // 连续发送失败计数
|
||
|
||
static int iot_mqtt_send(const uint8_t *buf, uint16_t len) {
|
||
uint16_t total_sent = 0;
|
||
int retries = 0;
|
||
while (total_sent < len) {
|
||
uint32_t slen = len - total_sent;
|
||
uint8_t ret = WCHNET_SocketSend(g_iot_socket, (uint8_t *)(buf + total_sent), &slen);
|
||
total_sent += (uint16_t)slen;
|
||
if (slen == 0 || ret != WCHNET_ERR_SUCCESS || ++retries > 10) {
|
||
_iot_send_fail_cnt++;
|
||
PRINT("IOT: SocketSend FAIL#%d ret=0x%02X sent=%u/%u\n",
|
||
_iot_send_fail_cnt, ret, total_sent, len);
|
||
if (_iot_send_fail_cnt >= 3) {
|
||
PRINT("IOT: SocketSend 3 consecutive failures, forcing reconnect\n");
|
||
WCHNET_SocketClose(g_iot_socket, TCP_CLOSE_NORMAL);
|
||
g_iot_socket = 0xFF;
|
||
g_iot_state = IOT_STATE_DISCONNECTED;
|
||
_iot_send_fail_cnt = 0;
|
||
_iot_reconnect_deadline = mstick() + 3000; // 3s 后重连, 让 WCHNET 清旧 socket
|
||
}
|
||
return -1;
|
||
}
|
||
}
|
||
_iot_send_fail_cnt = 0; // 成功一次就清零
|
||
return 0;
|
||
}
|
||
|
||
/* 发送 MQTT CONNECT — 参考 net_srv.c mqtt_connect 实现 */
|
||
static int iot_mqtt_send_connect(void) {
|
||
static MQTTPacket_connectData opts; // static: 避免栈开销
|
||
static uint8_t buf[256];
|
||
int len;
|
||
|
||
// 手动初始化(参考 net_srv.c 模式,不依赖 brace-initializer 赋值)
|
||
memset(&opts, 0, sizeof(opts));
|
||
opts.struct_id[0] = 'M'; opts.struct_id[1] = 'Q';
|
||
opts.struct_id[2] = 'T'; opts.struct_id[3] = 'C';
|
||
opts.MQTTVersion = 4;
|
||
opts.keepAliveInterval = IOT_MQTT_KEEPALIVE_SEC;
|
||
opts.cleansession = 1;
|
||
|
||
// clientID — 参考 net_srv.c 直接使用 g_dev_number_str
|
||
opts.clientID.cstring = g_dev_number_str;
|
||
|
||
// username / password — 参考 net_srv.c 直接赋值
|
||
opts.username.cstring = (char *)iot_net_info.username;
|
||
opts.password.cstring = (char *)iot_net_info.password;
|
||
|
||
len = MQTTSerialize_connect(buf, sizeof(buf), &opts);
|
||
if (len <= 0) {
|
||
PRINT("IOT: MQTTSerialize_connect failed\n");
|
||
return -1;
|
||
}
|
||
PRINT("IOT: → CONNECT clientId=%s keepAlive=%d\n",
|
||
opts.clientID.cstring, opts.keepAliveInterval);
|
||
return iot_mqtt_send(buf, (uint16_t)len);
|
||
}
|
||
|
||
/* 发送 MQTT SUBSCRIBE — V1.01 双主题协议 */
|
||
static int iot_mqtt_send_subscribe(void) {
|
||
static uint8_t buf[256]; // static: 避免栈溢出
|
||
int len;
|
||
char topic[IOT_MQTT_TOPIC_MAX_LEN];
|
||
MQTTString topics[1];
|
||
int qos[1] = {1};
|
||
int count = 0;
|
||
|
||
// V1.01: 仅订阅 dld960/{sn}/srv 双主题
|
||
snprintf(topic, sizeof(topic), "dld960/%s/srv", g_iot_dev_serial);
|
||
topics[count].cstring = topic; topics[count].lenstring.len = 0; count++;
|
||
|
||
len = MQTTSerialize_subscribe(buf, sizeof(buf), 0,
|
||
++g_iot_msg_id, count, topics, qos);
|
||
if (len <= 0) {
|
||
PRINT("IOT: MQTTSerialize_subscribe failed\n");
|
||
return -1;
|
||
}
|
||
PRINT("IOT: → SUBSCRIBE pktId=%d topics=%d\n", g_iot_msg_id, count);
|
||
return iot_mqtt_send(buf, (uint16_t)len);
|
||
}
|
||
|
||
/* 发送 MQTT PUBLISH */
|
||
static int iot_mqtt_publish(const char *topic, const char *payload,
|
||
uint16_t payload_len, uint8_t qos) {
|
||
static uint8_t buf[IOT_MQTT_SEND_BUF_LEN]; // static: 避免 2KB 栈溢出
|
||
int len;
|
||
MQTTString mqtt_topic = MQTTString_initializer;
|
||
mqtt_topic.cstring = (char *)topic;
|
||
|
||
len = MQTTSerialize_publish(buf, sizeof(buf), 0, qos, 0,
|
||
++g_iot_msg_id, mqtt_topic,
|
||
(unsigned char *)payload, payload_len);
|
||
if (len <= 0) {
|
||
PRINT("IOT: MQTTSerialize_publish failed\n");
|
||
return -1;
|
||
}
|
||
return iot_mqtt_send(buf, (uint16_t)len);
|
||
}
|
||
|
||
/*===========================================================================
|
||
* MQTT 接收处理
|
||
*===========================================================================*/
|
||
|
||
/* 处理一条收到的 MQTT PUBLISH 消息 */
|
||
static void iot_handle_publish(const char *topic, uint8_t *payload, int payload_len) {
|
||
static char json[IOT_MQTT_RECV_BUF_LEN]; // static: 避免 1KB 栈开销
|
||
int copy_len = payload_len < (int)sizeof(json) - 1 ? payload_len : (int)sizeof(json) - 1;
|
||
memcpy(json, payload, copy_len);
|
||
json[copy_len] = '\0';
|
||
|
||
PRINT("IOT: PUBLISH topic=%s payload=%s\n", topic, json);
|
||
|
||
/* 提取 msg_id 和 cmd */
|
||
char tmp[256];
|
||
uint32_t msg_id = 0;
|
||
|
||
// 简单提取: 解析 JSON 中的 msg_id 和 cmd
|
||
memset(tmp, 0, sizeof(tmp));
|
||
simple_parse_json(json, "\"msg_id\"", tmp);
|
||
if (strlen(tmp) > 0) msg_id = (uint32_t)strtoul(tmp, NULL, 10);
|
||
|
||
memset(tmp, 0, sizeof(tmp));
|
||
simple_parse_json(json, "\"cmd\"", tmp);
|
||
// strip quotes
|
||
char *cmd_str = tmp;
|
||
if (cmd_str[0] == '"') cmd_str++;
|
||
int cmd_len = strlen(cmd_str);
|
||
if (cmd_len > 0 && cmd_str[cmd_len - 1] == '"') cmd_str[cmd_len - 1] = '\0';
|
||
|
||
if (strlen(cmd_str) == 0) {
|
||
PRINT("IOT: no cmd in PUBLISH\n");
|
||
return;
|
||
}
|
||
|
||
/* 构建响应 topic — V1.01 双主题 */
|
||
char resp_topic[IOT_MQTT_TOPIC_MAX_LEN];
|
||
snprintf(resp_topic, sizeof(resp_topic), "dld960/%s/dev", g_iot_dev_serial);
|
||
|
||
/* 处理命令 — 复用 TCP JSON 的命令逻辑 */
|
||
/* 目前仅实现基本响应框架, 后续逐步对接 Loop MCU 命令 */
|
||
|
||
if (strcmp(cmd_str, "dev_info_query") == 0) {
|
||
char data_json[512];
|
||
snprintf(data_json, sizeof(data_json),
|
||
"{\"msg_id\":%lu,\"cmd\":\"dev_info_query\","
|
||
"\"ts\":%lu,\"code\":0,\"msg\":\"success\","
|
||
"\"data\":{"
|
||
"\"dev_serial\":\"%s\",\"hard_ver\":\"%s\",\"soft_ver\":\"%s\","
|
||
"\"model\":\"%s\",\"product_code\":\"960001\","
|
||
"\"sub_code\":{\"net\":%s,\"iot\":%s},"
|
||
"\"bus\":{\"bus1\":0,\"bus2\":0,\"bus3\":0,\"bus4\":0}}}",
|
||
msg_id, dev_time_now(),
|
||
g_dev_number_str, HARDWARE_VER, FIRMWARE_VER, PRODUCT_MODEL,
|
||
g_sub_code_enable.net_enable ? "true" : "false",
|
||
g_sub_code_enable.iot_enable ? "true" : "false");
|
||
iot_mqtt_publish(resp_topic, data_json, strlen(data_json), 1);
|
||
|
||
} else if (strcmp(cmd_str, "report_config") == 0) {
|
||
/* 平台下发配置 → 更新 g_report_cfg, 回 ACK */
|
||
char data[512] = {0};
|
||
simple_parse_json(json, "\"data\"", data);
|
||
if (strlen(data) > 0) {
|
||
char tmp2[64];
|
||
memset(tmp2, 0, sizeof(tmp2));
|
||
simple_parse_json(data, "\"enable\"", tmp2);
|
||
if (strlen(tmp2) > 0) g_report_cfg.enable = (strcmp(tmp2, "true") == 0);
|
||
memset(tmp2, 0, sizeof(tmp2));
|
||
simple_parse_json(data, "\"interval\"", tmp2);
|
||
if (strlen(tmp2) > 0) g_report_cfg.interval = (uint16_t)strtoul(tmp2, NULL, 10);
|
||
memset(tmp2, 0, sizeof(tmp2));
|
||
simple_parse_json(data, "\"once\"", tmp2);
|
||
if (strlen(tmp2) > 0) g_report_cfg.once = (strcmp(tmp2, "true") == 0);
|
||
memset(tmp2, 0, sizeof(tmp2));
|
||
simple_parse_json(data, "\"env_eval\"", tmp2);
|
||
if (strlen(tmp2) > 0) g_report_cfg.env_eval = (strcmp(tmp2, "true") == 0);
|
||
memset(tmp2, 0, sizeof(tmp2));
|
||
simple_parse_json(data, "\"sensor_type\"", tmp2);
|
||
if (strlen(tmp2) > 0) g_report_cfg.sensor_type = (uint8_t)strtoul(tmp2, NULL, 10);
|
||
}
|
||
char resp[256];
|
||
snprintf(resp, sizeof(resp),
|
||
"{\"msg_id\":%lu,\"cmd\":\"report_config\","
|
||
"\"ts\":%lu,\"code\":0,\"msg\":\"success\"}",
|
||
msg_id, dev_time_now());
|
||
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
|
||
|
||
} else if (strcmp(cmd_str, "event_report") == 0) {
|
||
/* 平台对 event_report 的 ACK: 检查 code 字段 */
|
||
char tmp3[64] = {0};
|
||
int code = -1;
|
||
simple_parse_json(json, "\"code\"", tmp3);
|
||
if (strlen(tmp3) > 0) code = (int)strtol(tmp3, NULL, 10);
|
||
if (code == 0) {
|
||
iot_evt_handle_ack(msg_id, 0);
|
||
} else {
|
||
iot_evt_handle_ack(msg_id, code);
|
||
}
|
||
// ACK 不回复 (避免乒乓)
|
||
|
||
} else if (strcmp(cmd_str, "pwd_verify") == 0) {
|
||
// 验证密码
|
||
char password[16] = {0};
|
||
char *data = (char *)malloc(512);
|
||
if (data) {
|
||
memset(data, 0, 512);
|
||
simple_parse_json(json, "\"data\"", data);
|
||
if (strlen(data) > 0) {
|
||
memset(tmp, 0, sizeof(tmp));
|
||
simple_parse_json(data, "\"password\"", tmp);
|
||
// strip quotes
|
||
char *pwd = tmp;
|
||
if (pwd[0] == '"') pwd++;
|
||
int plen = strlen(pwd);
|
||
if (plen > 0 && pwd[plen - 1] == '"') pwd[plen - 1] = '\0';
|
||
strncpy(password, pwd, 15);
|
||
}
|
||
free(data);
|
||
}
|
||
|
||
if (strlen(password) == 6 && memcmp(password, g_dev_password, 6) == 0) {
|
||
char resp[256];
|
||
snprintf(resp, sizeof(resp),
|
||
"{\"msg_id\":%lu,\"cmd\":\"pwd_verify\","
|
||
"\"ts\":%lu,\"code\":0,\"msg\":\"success\"}",
|
||
msg_id, dev_time_now());
|
||
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
|
||
PRINT("IOT: Auth success via MQTT\n");
|
||
} else {
|
||
char resp[256];
|
||
snprintf(resp, sizeof(resp),
|
||
"{\"msg_id\":%lu,\"cmd\":\"pwd_verify\","
|
||
"\"ts\":%lu,\"code\":2,\"msg\":\"password incorrect\"}",
|
||
msg_id, dev_time_now());
|
||
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
|
||
}
|
||
|
||
} else {
|
||
// 暂不支持的命令, 返回错误
|
||
char resp[256];
|
||
snprintf(resp, sizeof(resp),
|
||
"{\"msg_id\":%lu,\"cmd\":\"%s\","
|
||
"\"ts\":%lu,\"code\":4,\"msg\":\"unsupported command\"}",
|
||
msg_id, cmd_str, dev_time_now());
|
||
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
|
||
}
|
||
}
|
||
|
||
/* 发送 initialize 上线消息 (IoT 路径) */
|
||
static void iot_send_initialize(void) {
|
||
char payload[512];
|
||
char topic[IOT_MQTT_TOPIC_MAX_LEN];
|
||
snprintf(payload, sizeof(payload),
|
||
"{\"msg_id\":%lu,\"cmd\":\"initialize\",\"ts\":%lu,"
|
||
"\"data\":{\"dev_serial\":\"%s\",\"model\":\"%s\","
|
||
"\"hard_ver\":\"%s\",\"soft_ver\":\"%s\","
|
||
"\"extra_info\":{\"code\":\"0\",\"csq\":\"0\",\"location\":\"\"}}}",
|
||
(unsigned long)(g_iot_msg_id + 1), dev_time_now(),
|
||
g_iot_dev_serial, PRODUCT_MODEL, HARDWARE_VER, FIRMWARE_VER);
|
||
snprintf(topic, sizeof(topic), "dld960/%s/dev", g_iot_dev_serial);
|
||
iot_mqtt_publish(topic, payload, strlen(payload), 0);
|
||
}
|
||
|
||
/* 处理 MQTT SUBACK */
|
||
static void iot_handle_suback(void) {
|
||
PRINT("IOT: ← SUBACK, topics subscribed\n");
|
||
g_iot_state = IOT_STATE_READY;
|
||
_iot_reconnect_backoff = 0; // 连接成功, 重置退避
|
||
|
||
// V1.03: 订阅成功后发布 initialize 告知服务器上线 (IoT 路径)
|
||
iot_send_initialize();
|
||
}
|
||
|
||
/* 处理收到的 MQTT 数据 */
|
||
static void iot_process_recv(void) {
|
||
uint8_t header;
|
||
int qos, payload_len;
|
||
unsigned char retained;
|
||
unsigned short packet_id;
|
||
MQTTString topic;
|
||
unsigned char *payload_ptr;
|
||
|
||
while (_iot_recv_len >= 2) {
|
||
header = _iot_recv_buf[0];
|
||
int msg_type = (header >> 4) & 0x0F;
|
||
|
||
/* 先尝试解析整个 MQTT 包 */
|
||
int mqtt_pkt_len = 0;
|
||
int multiplier = 1;
|
||
int rem_len = 0;
|
||
int i = 1;
|
||
while (i < (int)_iot_recv_len && i < 5) {
|
||
rem_len += (_iot_recv_buf[i] & 0x7F) * multiplier;
|
||
multiplier *= 128;
|
||
if ((_iot_recv_buf[i] & 0x80) == 0) {
|
||
mqtt_pkt_len = 1 + (i - 1 + 1) + rem_len; // header + rem_len_bytes + payload
|
||
break;
|
||
}
|
||
i++;
|
||
}
|
||
if (mqtt_pkt_len == 0 || mqtt_pkt_len > (int)_iot_recv_len) {
|
||
return; // 包不完整, 等更多数据
|
||
}
|
||
|
||
PRINT("IOT: ← MQTT packet type=%d len=%d\n", msg_type, mqtt_pkt_len);
|
||
|
||
switch (msg_type) {
|
||
case CONNACK: {
|
||
unsigned char session_present, connack_rc;
|
||
if (MQTTDeserialize_connack(&session_present, &connack_rc,
|
||
_iot_recv_buf, mqtt_pkt_len) == 1) {
|
||
PRINT("IOT: ← CONNACK rc=%d\n", connack_rc);
|
||
if (connack_rc == 0) {
|
||
// 连接成功 → 订阅
|
||
g_iot_state = IOT_STATE_MQTT_CONNECTED;
|
||
iot_mqtt_send_subscribe();
|
||
} else {
|
||
PRINT("IOT: CONNACK rejected, rc=%d\n", connack_rc);
|
||
g_iot_state = IOT_STATE_DISCONNECTED;
|
||
}
|
||
}
|
||
break;
|
||
}
|
||
|
||
case SUBACK:
|
||
iot_handle_suback();
|
||
break;
|
||
|
||
case PUBLISH: {
|
||
unsigned char pub_dup;
|
||
if (MQTTDeserialize_publish(&pub_dup, &qos, &retained,
|
||
&packet_id, &topic,
|
||
&payload_ptr, &payload_len,
|
||
_iot_recv_buf, mqtt_pkt_len) == 1) {
|
||
char topic_str[IOT_MQTT_TOPIC_MAX_LEN];
|
||
memcpy(topic_str, topic.cstring ? topic.cstring : "",
|
||
topic.lenstring.len < (int)sizeof(topic_str) - 1
|
||
? topic.lenstring.len : (int)sizeof(topic_str) - 1);
|
||
topic_str[topic.lenstring.len] = '\0';
|
||
|
||
iot_handle_publish(topic_str, payload_ptr, payload_len);
|
||
|
||
// 如果 QoS > 0, 发送 PUBACK
|
||
if (qos > 0) {
|
||
uint8_t ack_buf[4];
|
||
int ack_len = MQTTSerialize_ack(ack_buf, sizeof(ack_buf),
|
||
PUBACK, 0, packet_id);
|
||
iot_mqtt_send(ack_buf, (uint16_t)ack_len);
|
||
}
|
||
}
|
||
break;
|
||
}
|
||
|
||
case PINGRESP:
|
||
PRINT("IOT: ← PINGRESP\n");
|
||
break;
|
||
|
||
default:
|
||
break;
|
||
}
|
||
|
||
/* 移除已处理的包 */
|
||
memmove(_iot_recv_buf, _iot_recv_buf + mqtt_pkt_len,
|
||
_iot_recv_len - mqtt_pkt_len);
|
||
_iot_recv_len -= mqtt_pkt_len;
|
||
}
|
||
}
|
||
|
||
/*===========================================================================
|
||
* 传感器数据上报 (三档调度)
|
||
* 立即档: car_state 翻转沿 → 直通门控立即上报 (进/出车不等间隔)
|
||
* 快速档: |variation| > 阈值 → 300ms
|
||
* 空闲档: 按配置 interval (默认 60s, 最小 1s)
|
||
*
|
||
* Step 1 — 消费 Loop MCU 新帧 → 事件沿检测 + 更新缓存
|
||
* Step 2 — 判定上报档位 (car_edge / fast_mode / idle)
|
||
* Step 3 — 间隔门控 (car_edge 直通)
|
||
*===========================================================================*/
|
||
void iot_mqtt_publish_sensor(void) {
|
||
static uint8_t _last_car_state[4] = {0};
|
||
|
||
/*--- Step 1: 消费 0xC0 帧 (Step1 直读路径, 竞争窗口小, 少数帧走此路) ---*/
|
||
/* 注意: 置于 READY/enable 检查之前 —— 断网/未使能期间事件照样入队,
|
||
重连后由 iot_evt_process() 补报 (协议 V1.04)。
|
||
多数帧由 uart_srv→lup 回调路径摄取, 两路都汇入 iot_sensor_ingest */
|
||
if (g_pkg_uart_2.flag != 0
|
||
&& g_pkg_uart_2.pkg[0] == 0x7F
|
||
&& g_pkg_uart_2.pkg[3] == 0xC0) {
|
||
|
||
LUP_SensorReport sr;
|
||
int ret = -100;
|
||
memset(&sr, 0, sizeof(sr));
|
||
/* 直读路径未经过 lup_process_frame, 须自行校验 checksum */
|
||
if (lup_verify_checksum(g_pkg_uart_2.pkg, g_pkg_uart_2.offset) == 0) {
|
||
ret = lup_parse_sensor_report(g_pkg_uart_2.pkg, g_pkg_uart_2.offset, &sr);
|
||
}
|
||
InitPkgUart(&g_pkg_uart_2); // 读完即清, 不持有
|
||
|
||
if (ret == 0) {
|
||
iot_sensor_ingest(&sr); // 统一摄取: 事件沿 + loop_data 缓存
|
||
} else {
|
||
PRINT("IOT: sensor frame invalid (%d)\n", ret);
|
||
}
|
||
}
|
||
|
||
/*--- 事件上报状态机: 发送/重发/重连补报 ---*/
|
||
iot_evt_process();
|
||
|
||
if (g_iot_state != IOT_STATE_READY) return;
|
||
if (!g_report_cfg.enable) return;
|
||
|
||
/* 尚无缓存数据 → 不发 */
|
||
if (!_cached_sr_valid) return;
|
||
|
||
/*--- Step 2: 判定上报档位 (三档) ---
|
||
car_edge : 任一通道 car_state 翻转沿 → 【立即上报】直通间隔门控
|
||
fast_mode : 任一通道 |variation|>阈值 → 300ms 快速档
|
||
否则 : 空闲档, 按配置 interval (默认 60s, 最小 1s) */
|
||
uint32_t interval_ms;
|
||
uint8_t fast_mode = 0;
|
||
uint8_t car_edge = 0;
|
||
|
||
{
|
||
uint8_t i;
|
||
for (i = 0; i < _cached_sr.coil_count; i++) {
|
||
/* car_state 翻转沿(进车/出车)优先级最高, 命中即定档 */
|
||
if (_last_car_state[i] != _cached_sr.coils[i].car_state) {
|
||
car_edge = 1;
|
||
break;
|
||
}
|
||
// variation 为有符号(V1.05): 正=车/裕量, 负=反向漂移, 两向显著变化都加速
|
||
int32_t v = _cached_sr.coils[i].variation;
|
||
int32_t av = (v >= 0) ? v : -v;
|
||
if (av > IOT_MQTT_VARIATION_THRESHOLD) {
|
||
fast_mode = 1; /* 记录但不 break: 后续通道可能有沿(优先级更高) */
|
||
}
|
||
}
|
||
}
|
||
|
||
if (fast_mode || car_edge) {
|
||
interval_ms = IOT_MQTT_FAST_INTERVAL_MS; // 300ms (car_edge 时门控被直通, 仅作记录)
|
||
} else {
|
||
/* 空闲态: 使用配置的 interval (秒→毫秒), 最小 1s */
|
||
uint32_t cfg_ms = (uint32_t)g_report_cfg.interval * 1000;
|
||
interval_ms = (cfg_ms >= IOT_MQTT_IDLE_INTERVAL_MIN_MS)
|
||
? cfg_ms : IOT_MQTT_IDLE_INTERVAL_MIN_MS;
|
||
}
|
||
|
||
/*--- Step 3: 间隔门控 (car_edge 直通: 进/出车立即上报, 不等间隔) ---
|
||
防洪泛天然有界: Loop MCU 事件帧最快 150ms 一帧, 且发布后快照即同步,
|
||
同一个沿不会重复触发立即档 */
|
||
{
|
||
uint32_t now = mstick();
|
||
if (!car_edge
|
||
&& _last_publish_ms != 0
|
||
&& (now - _last_publish_ms) < interval_ms) {
|
||
return; // 未到间隔; _last_car_state 不刷新, 翻转沿保持"待发"
|
||
}
|
||
_last_publish_ms = now;
|
||
}
|
||
|
||
/* 确定发布 → 刷新 car_state 快照。必须放在门控之后:
|
||
若放门控前, 沿被门控吞掉时快照已刷新, 下一轮检测不到翻转,
|
||
事件退化为空闲间隔上报, 快速上报失效 */
|
||
{
|
||
uint8_t i;
|
||
for (i = 0; i < _cached_sr.coil_count; i++) {
|
||
_last_car_state[i] = _cached_sr.coils[i].car_state;
|
||
}
|
||
}
|
||
|
||
/*--- Step 4: 构建 JSON → 发布 (字段与 V1.02 协议一致) ---
|
||
直接构建到 payload, 省掉 data_json[1024]:
|
||
- 用 iot_mqtt_publish() 代替 mqtt_publish(), 发送缓冲隔离
|
||
(event_report→mqttBuf vs loop_data→iot_mqtt_publish::buf[1024])
|
||
- coil_count 边界硬限, 防 0xC0 坏帧致溢出 */
|
||
static char payload[1400];
|
||
char *p = payload;
|
||
int remaining = sizeof(payload);
|
||
int written;
|
||
const char *freq_level_names[] = {"high", "mid_high", "mid_low", "low"};
|
||
uint8_t i;
|
||
uint8_t coil_n = _cached_sr.coil_count;
|
||
uint32_t my_msg_id = g_iot_msg_id + 1; /* iot_mqtt_publish 内部 ++g_iot_msg_id, 这里只预读 */
|
||
if (coil_n > 4) coil_n = 4; /* 硬限: 防坏帧致 snprintf 循环溢出 */
|
||
|
||
/* 先写外层包装头 (msg_id/cmd/ts/data) */
|
||
written = snprintf(p, remaining,
|
||
"{\"msg_id\":%lu,\"cmd\":\"loop_data\",\"ts\":%lu,\"data\":{\"channels\":[",
|
||
(unsigned long)my_msg_id, dev_time_now());
|
||
if (written < 0 || written >= remaining) return;
|
||
p += written; remaining -= written;
|
||
|
||
for (i = 0; i < coil_n; i++) {
|
||
const LUP_CoilSensor *cs = &_cached_sr.coils[i];
|
||
const char *misc_type_str = "time";
|
||
uint32_t misc_val = 0;
|
||
|
||
if (cs->misc_type == 0) { misc_type_str = "time"; misc_val = cs->misc.passtime_ms; }
|
||
else if (cs->misc_type == 1) { misc_type_str = "cut_count"; misc_val = cs->misc.cut_amount; }
|
||
else if (cs->misc_type == 2) { misc_type_str = "flow_count"; misc_val = cs->misc.flow_amount; }
|
||
else if (cs->misc_type == 3) { misc_type_str = "relay_count"; misc_val = cs->misc.relay_count; }
|
||
|
||
written = snprintf(p, remaining,
|
||
"%s{\"ch\":%d,\"level\":\"%s\",\"iscar\":%s,"
|
||
"\"loop_ok\":%s,\"freq\":%lu,\"diff\":%d,"
|
||
"\"sens\":%d,\"cndtn\":%d,"
|
||
"\"misc\":{\"type\":\"%s\",\"value\":%lu}}",
|
||
(i > 0) ? "," : "",
|
||
i + 1,
|
||
freq_level_names[cs->freq_level],
|
||
cs->car_state ? "true" : "false",
|
||
cs->loop_state ? "false" : "true",
|
||
cs->freq, cs->variation,
|
||
cs->sensitivity, cs->condition,
|
||
misc_type_str, misc_val);
|
||
if (written < 0 || written >= remaining) return;
|
||
p += written; remaining -= written;
|
||
}
|
||
|
||
snprintf(p, remaining, "]}}");
|
||
|
||
/* 用 iot_mqtt_publish: 发送缓冲 buf[1024] 与 mqttBuf[1024] 物理隔离,
|
||
避免 event_report 的 MQTT 二进制污染 loop_data 载荷 (2026-07-23 _raw 事故 v3) */
|
||
{
|
||
char topic[IOT_MQTT_TOPIC_MAX_LEN];
|
||
snprintf(topic, sizeof(topic), "dld960/%s/dev", g_iot_dev_serial);
|
||
iot_mqtt_publish(topic, payload, strlen(payload), 0);
|
||
}
|
||
}
|
||
|
||
/*===========================================================================
|
||
* 连接管理
|
||
*===========================================================================*/
|
||
|
||
/* 连接到 MQTT Broker — 重新创建 TCP socket 并发起连接 */
|
||
static void iot_connect_broker(void) {
|
||
WCHNET_CreateTcpMqttSocket(); // 重新 create + connect (不是复用旧的!)
|
||
g_iot_socket = SocketId_TCP;
|
||
g_iot_state = IOT_STATE_TCP_CONNECTING;
|
||
_iot_connect_start = mstick();
|
||
PRINT("IOT: TCP connecting to broker (sock=%d)...\n", g_iot_socket);
|
||
}
|
||
|
||
/*===========================================================================
|
||
* Socket 中断处理
|
||
*===========================================================================*/
|
||
void iot_mqtt_handle_sock_int(uint8_t socketid, uint8_t intstat) {
|
||
if (socketid != g_iot_socket) return;
|
||
|
||
/* CONNECT 成功 — 在中断上下文直接发 MQTT CONNECT
|
||
(对齐旧 mqtt_connect 行为: WCHNET SocketSend 在中断外会 hard fault) */
|
||
if (intstat & SINT_STAT_CONNECT) {
|
||
PRINT("IOT: TCP connected (sock=%d)\n", socketid);
|
||
g_iot_state = IOT_STATE_TCP_CONNECTED;
|
||
_iot_recv_len = 0;
|
||
_iot_send_fail_cnt = 0; // 新连接, 重置失败计数
|
||
iot_mqtt_send_connect(); // 在中断上下文发, 不延后到 poll
|
||
g_iot_state = IOT_STATE_MQTT_CONNECTING;
|
||
}
|
||
|
||
/* 收到数据 */
|
||
if (intstat & SINT_STAT_RECV) {
|
||
uint32_t recv_len = WCHNET_SocketRecvLen(socketid, NULL);
|
||
if (recv_len > 0) {
|
||
uint16_t space = IOT_MQTT_RECV_BUF_LEN - _iot_recv_len;
|
||
if (recv_len > space) recv_len = space;
|
||
uint32_t rd_len = recv_len;
|
||
static uint8_t tmp[RECE_BUF_LEN]; /* static: 1152B 栈分配在中断上下文中会撑爆 2KB 栈 (2026-07-23) */
|
||
WCHNET_SocketRecv(socketid, tmp, &rd_len);
|
||
memcpy(_iot_recv_buf + _iot_recv_len, tmp, (uint16_t)rd_len);
|
||
_iot_recv_len += (uint16_t)rd_len;
|
||
iot_process_recv();
|
||
}
|
||
}
|
||
|
||
/* 断开 / 超时 */
|
||
if (intstat & (SINT_STAT_DISCONNECT | SINT_STAT_TIM_OUT)) {
|
||
PRINT("IOT: socket disconnect/timeout (sock=%d)\n", socketid);
|
||
g_iot_state = IOT_STATE_DISCONNECTED;
|
||
g_iot_socket = 0xFF;
|
||
_iot_recv_len = 0;
|
||
_cached_sr_valid = 0; // 清缓存: 重连后等新帧
|
||
_last_publish_ms = 0; // 复位: 重连后立即首发
|
||
_iot_reconnect_backoff = IOT_MQTT_RECONNECT_MIN_MS;
|
||
_iot_reconnect_deadline = mstick() + _iot_reconnect_backoff;
|
||
}
|
||
}
|
||
|
||
/*===========================================================================
|
||
* 硬件看门狗 IWDG — 主循环卡死保护
|
||
* LSI≈40kHz, Prescaler=256, Reload=625 → ~4s 超时
|
||
* 保活链路: MQTT PINGREQ/PINGRESP → TCP KeepAlive → IWDG
|
||
* (heartbeat 是单向设备上报, 平台不回复, 不能用作存活检测)
|
||
*===========================================================================*/
|
||
|
||
/* IWDG 初始化 */
|
||
void iot_watchdog_init(void) {
|
||
IWDG_WriteAccessCmd(IWDG_WriteAccess_Enable);
|
||
IWDG_SetPrescaler(IWDG_Prescaler_256);
|
||
IWDG_SetReload(625); // 625 × (256/40000) ≈ 4.0s
|
||
IWDG_ReloadCounter();
|
||
IWDG_Enable();
|
||
PRINT("IOT: IWDG init ok, timeout≈4s\n");
|
||
}
|
||
|
||
/* 喂狗 — 主循环每轮调用 */
|
||
void iot_watchdog_kick(void) {
|
||
IWDG_ReloadCounter();
|
||
}
|
||
|
||
/*===========================================================================
|
||
* 主循环轮询
|
||
*===========================================================================*/
|
||
void iot_mqtt_poll(void) {
|
||
iot_watchdog_kick(); // 喂硬件 IWDG
|
||
|
||
/* 断线重连 */
|
||
if (g_iot_state == IOT_STATE_DISCONNECTED && g_iot_socket == 0xFF) {
|
||
if (_iot_reconnect_deadline == 0 || mstick() > _iot_reconnect_deadline) {
|
||
iot_connect_broker();
|
||
if (g_iot_socket == 0xFF) {
|
||
// 连接失败, 退避
|
||
if (_iot_reconnect_backoff == 0) {
|
||
_iot_reconnect_backoff = IOT_MQTT_RECONNECT_MIN_MS;
|
||
} else {
|
||
_iot_reconnect_backoff *= 2;
|
||
if (_iot_reconnect_backoff > IOT_MQTT_RECONNECT_MAX_MS)
|
||
_iot_reconnect_backoff = IOT_MQTT_RECONNECT_MAX_MS;
|
||
}
|
||
_iot_reconnect_deadline = mstick() + _iot_reconnect_backoff;
|
||
PRINT("IOT: reconnect in %lu ms\n", _iot_reconnect_backoff);
|
||
}
|
||
}
|
||
return;
|
||
}
|
||
|
||
/* TCP 连接超时 (10s) */
|
||
if (g_iot_state == IOT_STATE_TCP_CONNECTING) {
|
||
if (mstick() - _iot_connect_start > 10000) {
|
||
PRINT("IOT: TCP connect timeout\n");
|
||
WCHNET_SocketClose(g_iot_socket, TCP_CLOSE_NORMAL);
|
||
g_iot_socket = 0xFF;
|
||
g_iot_state = IOT_STATE_DISCONNECTED;
|
||
return;
|
||
}
|
||
}
|
||
|
||
/* TCP 已连接 — 发送 MQTT CONNECT(延迟到 poll 而非中断内) */
|
||
if (g_iot_state == IOT_STATE_TCP_CONNECTED) {
|
||
iot_mqtt_send_connect();
|
||
g_iot_state = IOT_STATE_MQTT_CONNECTING;
|
||
}
|
||
|
||
/* MQTT Keepalive — PINGREQ (60s, broker 不回 → TCP 超时 → 重连)
|
||
loop_data 按平台下发 g_report_cfg.interval 上报, 无需额外 heartbeat JSON */
|
||
if (g_iot_state == IOT_STATE_READY) {
|
||
uint32_t now = mstick();
|
||
if (now - _iot_last_heartbeat > IOT_MQTT_HEARTBEAT_MS) {
|
||
uint8_t ping_buf[2];
|
||
int ping_len = MQTTSerialize_pingreq(ping_buf, sizeof(ping_buf));
|
||
iot_mqtt_send(ping_buf, (uint16_t)ping_len);
|
||
_iot_last_heartbeat = now;
|
||
}
|
||
}
|
||
}
|
||
|
||
/*===========================================================================
|
||
* 初始化
|
||
*===========================================================================*/
|
||
void iot_mqtt_init(void) {
|
||
/* 复用 net_srv.c 创建的 MQTT socket (SocketId_TCP) */
|
||
g_iot_socket = SocketId_TCP;
|
||
|
||
/* 初始化设备序列码 (12 hex chars) */
|
||
snprintf(g_iot_dev_serial, sizeof(g_iot_dev_serial),
|
||
"%02X%02X%02X%02X%02X%02X",
|
||
g_dev_number[0], g_dev_number[1], g_dev_number[2],
|
||
g_dev_number[3], g_dev_number[4], g_dev_number[5]);
|
||
|
||
PRINT("IOT: dev_serial=%s\n", g_iot_dev_serial);
|
||
|
||
/* 初始化传感器回调 — uart_srv 消费路径的 0xC0 帧喂入事件沿检测 (V1.04) */
|
||
lup_set_sensor_callback(iot_evt_sensor_cb);
|
||
|
||
_iot_reconnect_backoff = 0;
|
||
_iot_reconnect_deadline = mstick() + 2000; // 启动后 2s 开始首次连接
|
||
g_iot_state = IOT_STATE_DISCONNECTED;
|
||
|
||
/* 初始化硬件看门狗 IWDG (LSI≈40kHz, 4s 超时, 主循环喂狗) */
|
||
iot_watchdog_init();
|
||
}
|