Files
vd_960/vd960DBN/BLE/OnlyUpdateApp_Peripheral/APP/iot_mqtt_srv.c
T
wangfq 5a1893cd1c feat(vd960DBN)+fix(DBNMQTTool): log_query 改 hex 原始字节上报 (2026-08-18)
背景: MQTT 快照流实测 MQTTSerialize_publish failed — JSON 化快照记录 ~810B/条 超 800B 发送缓冲
方案(用户拍板): 对齐 BLE 通道, 原始字节 hex 上报

固件 (V4.3):
- offlog.c/h: 新增 offlog_evt_to_hex() (32B→64 hex)
- snapshot.c/h: 新增 snap_rec_to_hex() (64B→128 hex); 删 SNAP_MAX_QUERY_JSON, 恢复 count=2
- tcp_json_srv.c / iot_mqtt_srv.c: log_query 改 {"seq":N,"hex":"..."}; SEND_BUF 保持 800
- 2 条快照 hex 响应 406B < 800B

文档: TCP JSON V1.03 / MQTT V1.07 §4.17 records 改 hex + 解析表引用 BLE §6.4/§7

工具: parse_offlog_hex/parse_snap_hex/offlog_payload_desc + hex 展示; 验证: gcc 9 断言 + 工具解析全过 + offscreen UI
2026-08-18 14:04:26 +08:00

1200 lines
54 KiB
C
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
******************************************************************************
* @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 "offlog.h"
#include "snapshot.h"
#include "fault_diag.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;
static IotMqttState _prev_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++;
/* 事件日志: 线圈事件 (主循环上下文) */
offlog_coil((uint8_t)type + 1, ch, value); /* sub: 1=进 2=出 3=断开 4=恢复 */
}
_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;
FAULT_MARKER(MK_EVT_FEED_IN);
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;
FAULT_MARKER(MK_EVT_FEED_OUT);
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;
}
FAULT_MARKER(MK_EVT_FEED_OUT);
}
/* 统一帧摄取: 事件沿检测 + 刷新 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) {
FAULT_MARKER(MK_INGEST_IN);
iot_evt_feed(sr); // 事件沿检测
memcpy(&_cached_sr, sr, sizeof(LUP_SensorReport)); // 刷新 loop_data 快照
_cached_sr_valid = 1;
snap_enqueue(sr); /* 快照入队: 中断安全, 主循环 snap_flush 落盘 */
FAULT_MARKER(MK_INGEST_OUT);
}
/* lup 传感回调: uart_srv 消费路径的帧从这里喂入 (lup_process_frame 已过校验) */
static void iot_evt_sensor_cb(const uint8_t *pkg, uint16_t len) {
LUP_SensorReport sr;
FAULT_MARKER(MK_EVT_CB_IN);
memset(&sr, 0, sizeof(sr));
if (lup_parse_sensor_report(pkg, len, &sr) == 0) {
iot_sensor_ingest(&sr);
}
FAULT_MARKER(MK_EVT_CB_OUT);
}
/* 序列化并发布队列头 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);
offlog_evt_giveup(_evt_pend_id); // 事件日志: 重试耗尽
_evt_pend_id = 0;
_evt_gaveup = 1; // 挂起: 新事件或重连时再触发
return;
}
_evt_retry++;
offlog_evt_retry(_evt_pend_id, _evt_retry); // 事件日志: ACK 超时重发
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 if (strcmp(cmd_str, "log_stat") == 0) {
/* 日志统计/分页定位: event/snapshot 流 (协议 V1.07) */
char stream_buf[16];
memset(stream_buf, 0, sizeof(stream_buf));
simple_parse_json(json, "\"stream\"", stream_buf);
char resp[300];
if (strcmp(stream_buf, "snapshot") == 0) {
uint32_t total = snap_count();
uint32_t seq_last = snap_seq_last();
uint32_t seq_first = (total > 0) ? (seq_last - total + 1) : 0;
snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"log_stat\","
"\"ts\":%lu,\"code\":0,\"msg\":\"success\","
"\"data\":{\"stream\":\"snapshot\",\"enabled\":%s,\"boot_seq\":%lu,"
"\"count\":%lu,\"capacity\":%lu,\"seq_first\":%lu,\"seq_last\":%lu}}",
msg_id, dev_time_now(),
snap_enabled() ? "true" : "false",
(unsigned long)snap_boot_seq(),
(unsigned long)total, (unsigned long)SNAP_MAX_RECORDS,
(unsigned long)seq_first, (unsigned long)seq_last);
PRINT("IOT: log_stat(snapshot) count=%lu seq_first=%lu seq_last=%lu\n",
(unsigned long)total, (unsigned long)seq_first, (unsigned long)seq_last);
} else {
uint32_t total = offlog_count();
uint32_t seq_last = offlog_seq_last();
uint32_t seq_first = (total > 0) ? (seq_last - total + 1) : 0;
snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"log_stat\","
"\"ts\":%lu,\"code\":0,\"msg\":\"success\","
"\"data\":{\"stream\":\"event\",\"enabled\":%s,\"boot_seq\":%lu,"
"\"count\":%lu,\"capacity\":%lu,\"seq_first\":%lu,\"seq_last\":%lu}}",
msg_id, dev_time_now(),
offlog_enabled() ? "true" : "false",
(unsigned long)offlog_boot_seq(),
(unsigned long)total, (unsigned long)OFFLOG_MAX_RECORDS,
(unsigned long)seq_first, (unsigned long)seq_last);
PRINT("IOT: log_stat count=%lu seq_first=%lu seq_last=%lu\n",
(unsigned long)total, (unsigned long)seq_first, (unsigned long)seq_last);
}
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
} else if (strcmp(cmd_str, "log_query") == 0) {
/* 分页拉取脱机日志: 按全局序号, event/snapshot 流, hex 原始字节上报 (协议 V1.07) */
static char resp[IOT_MQTT_SEND_BUF_LEN]; /* static: 避免大栈开销 */
char stream_buf[16];
uint32_t start_seq = 0, req_count = 0;
int is_snap = 0;
memset(tmp, 0, sizeof(tmp));
simple_parse_json(json, "\"start_seq\"", tmp);
if (strlen(tmp) > 0) start_seq = (uint32_t)strtoul(tmp, NULL, 10);
memset(tmp, 0, sizeof(tmp));
simple_parse_json(json, "\"count\"", tmp);
if (strlen(tmp) > 0) req_count = (uint32_t)strtoul(tmp, NULL, 10);
memset(stream_buf, 0, sizeof(stream_buf));
simple_parse_json(json, "\"stream\"", stream_buf);
if (strcmp(stream_buf, "snapshot") == 0) is_snap = 1;
int pos = snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"log_query\",\"ts\":%lu,\"code\":0,"
"\"msg\":\"success\",\"data\":{\"start_seq\":%lu,\"records\":[",
msg_id, dev_time_now(), (unsigned long)start_seq);
uint32_t fetched = 0;
if (is_snap) {
uint32_t total = snap_count();
uint32_t seq_last = snap_seq_last();
uint32_t seq_first = (total > 0) ? (seq_last - total + 1) : 0;
if (req_count > SNAP_MAX_QUERY_RECORDS) req_count = SNAP_MAX_QUERY_RECORDS;
if (req_count == 0) req_count = SNAP_MAX_QUERY_RECORDS;
if (snap_enabled() && total > 0 && start_seq >= seq_first && start_seq <= seq_last) {
uint32_t idx = start_seq - seq_first; /* 逻辑索引 = 全局序号 - seq_first */
while (fetched < req_count && (idx + fetched) < total) {
SnapRec rec;
char hex[132]; /* 64B -> 128 hex + NUL */
if (snap_read_idx(idx + fetched, &rec) != 0) break;
if (fetched > 0) pos += snprintf(resp + pos, sizeof(resp) - pos, ",");
snap_rec_to_hex(&rec, hex, sizeof(hex));
pos += snprintf(resp + pos, sizeof(resp) - pos,
"{\"seq\":%lu,\"hex\":\"%s\"}", (unsigned long)rec.seq, hex);
fetched++;
}
}
PRINT("IOT: log_query(snapshot) start_seq=%lu req=%lu fetched=%lu\n",
(unsigned long)start_seq, (unsigned long)req_count, (unsigned long)fetched);
} else {
uint32_t total = offlog_count();
uint32_t seq_last = offlog_seq_last();
uint32_t seq_first = (total > 0) ? (seq_last - total + 1) : 0;
if (req_count > OFFLOG_MAX_QUERY_RECORDS) req_count = OFFLOG_MAX_QUERY_RECORDS;
if (req_count == 0) req_count = OFFLOG_MAX_QUERY_RECORDS;
if (total > 0 && start_seq >= seq_first && start_seq <= seq_last) {
uint32_t idx = start_seq - seq_first; /* 逻辑索引 = 全局序号 - seq_first */
while (fetched < req_count && (idx + fetched) < total) {
OfflogEvt evt;
char hex[68]; /* 32B -> 64 hex + NUL */
if (offlog_read_idx((uint16_t)(idx + fetched), &evt) != 0) break;
if (fetched > 0) pos += snprintf(resp + pos, sizeof(resp) - pos, ",");
offlog_evt_to_hex(&evt, hex, sizeof(hex));
pos += snprintf(resp + pos, sizeof(resp) - pos,
"{\"seq\":%lu,\"hex\":\"%s\"}", (unsigned long)evt.seq, hex);
fetched++;
}
}
PRINT("IOT: log_query start_seq=%lu req=%lu fetched=%lu\n",
(unsigned long)start_seq, (unsigned long)req_count, (unsigned long)fetched);
}
snprintf(resp + pos, sizeof(resp) - pos, "]}}");
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
} else if (strcmp(cmd_str, "log_clear") == 0) {
/* 清除日志 (审计留痕). event 阻塞 ~2.8s (63 扇区), snapshot 阻塞 ~45ms (逻辑清除+当前扇区) */
char stream_buf[16];
memset(stream_buf, 0, sizeof(stream_buf));
simple_parse_json(json, "\"stream\"", stream_buf);
char resp[256];
if (strcmp(stream_buf, "snapshot") == 0) {
snap_clear();
PRINT("IOT: log_clear(snapshot) done\n");
} else {
offlog_clear();
PRINT("IOT: log_clear done\n");
}
snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"log_clear\","
"\"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, "ssc_net_query") == 0) {
/* SSC 网络配置查询 (只读全局, 协议 V1.02) */
char resp[400];
snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"ssc_net_query\","
"\"ts\":%lu,\"code\":0,\"msg\":\"success\","
"\"data\":{\"dev_ip\":\"%d.%d.%d.%d\",\"subnet_mask\":\"%d.%d.%d.%d\","
"\"route_ip\":\"%d.%d.%d.%d\",\"lssc_ip\":\"%d.%d.%d.%d\","
"\"dns\":\"%d.%d.%d.%d\",\"port\":%d}}",
msg_id, dev_time_now(),
local_net_cfg.lip[0], local_net_cfg.lip[1], local_net_cfg.lip[2], local_net_cfg.lip[3],
local_net_cfg.sub[0], local_net_cfg.sub[1], local_net_cfg.sub[2], local_net_cfg.sub[3],
local_net_cfg.gw[0], local_net_cfg.gw[1], local_net_cfg.gw[2], local_net_cfg.gw[3],
net_center_info.lssc_ip[0], net_center_info.lssc_ip[1],
net_center_info.lssc_ip[2], net_center_info.lssc_ip[3],
local_net_cfg.dns[0], local_net_cfg.dns[1], local_net_cfg.dns[2], local_net_cfg.dns[3],
net_center_info.tcp_port);
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
PRINT("IOT: ssc_net_query dev_ip=%d.%d.%d.%d\n",
local_net_cfg.lip[0], local_net_cfg.lip[1], local_net_cfg.lip[2], local_net_cfg.lip[3]);
} else if (strcmp(cmd_str, "iot_net_query") == 0) {
/* IoT 网络配置查询 (只读全局) */
char resp[400];
snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"iot_net_query\","
"\"ts\":%lu,\"code\":0,\"msg\":\"success\","
"\"data\":{\"host\":\"%s\",\"port\":%d,\"client_id\":\"%s\","
"\"username\":\"%s\",\"password\":\"%s\"}}",
msg_id, dev_time_now(),
iot_net_info.remote_addr, iot_net_info.mqtt_port,
iot_net_info.client_id, iot_net_info.username, iot_net_info.password);
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
PRINT("IOT: iot_net_query host=%s port=%d\n",
iot_net_info.remote_addr, iot_net_info.mqtt_port);
} else if (strcmp(cmd_str, "iot_topic_query") == 0) {
/* IoT Topic 配置查询 (只读全局) */
char resp[300];
snprintf(resp, sizeof(resp),
"{\"msg_id\":%lu,\"cmd\":\"iot_topic_query\","
"\"ts\":%lu,\"code\":0,\"msg\":\"success\","
"\"data\":{\"client_id_enable\":%s,\"topic_pub\":\"%s\",\"topic_sub\":\"%s\"}}",
msg_id, dev_time_now(),
g_iot_topic.clientid_enable ? "true" : "false",
g_iot_topic.topic_pub, g_iot_topic.topic_sub);
iot_mqtt_publish(resp_topic, resp, strlen(resp), 1);
PRINT("IOT: iot_topic_query pub=%s sub=%s\n",
g_iot_topic.topic_pub, g_iot_topic.topic_sub);
} 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);
offlog_iot_disconnect(3); // 事件日志: broker 拒绝
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[800]; /* 2026-08-17: 1400→800, loop_data 最大~604B, 省 600B .bss */
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
/* 事件日志: MQTT 状态沿检测 (主循环上下文, 不在 socket 中断里写 SPI)
真复位 → 会有 BOOT 事件 + boot_seq 递增; 仅断连重连 → 只有以下网络事件 */
if (g_iot_state != _prev_iot_state) {
switch (g_iot_state) {
case IOT_STATE_MQTT_CONNECTED: offlog_iot_connect(); break; // TCP+CONNACK 成功
case IOT_STATE_READY: offlog_iot_ready(); break; // SUBACK → 将发 initialize
case IOT_STATE_DISCONNECTED: offlog_iot_disconnect(1); break; // 断连(细分原因在各分支单独记)
default: break;
}
_prev_iot_state = g_iot_state;
}
/* 断线重连 */
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;
offlog_iot_reconn(_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");
offlog_iot_disconnect(4); // 事件日志: 连接超时
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();
}