feat(iot_mqtt): 实现 event_report 平台必答+设备重发 (协议 V1.04); 时间单位统一 50ms
时间单位: - 串口协议 V1.06 / TCP JSON 协议: 杂项时间量、事件 value 由 5ms 修正为 50ms, 对齐 Loop 固件 50ms tick 实现 (TMR15 5ms×10) event_report (iot_mqtt_srv.c + net_srv.c): - 沿检测双路汇聚: Step1 消费路径 + lup 回调覆盖 uart_srv 路径, 不漏帧 - 16 深环形队列, 事件仅 ACK(code=0) 后出队; 溢出丢最旧 - 5s 超时重发同 msg_id/原始 ts ×3 次, 耗尽挂起, 新事件/重连沿合并补报 - 帧消费提前至 READY/enable 之前: 断网期间事件入队, 重连补报 - 线圈断开期间屏蔽 car 沿; loop_restore value=本地计时/50 - msg_id 独立 uint32 计数 (g_iot_msg_id 为 uint8 混用会破坏去重窗口) - net_srv.c ACK 路由: cmd=event_report 回显帧确认出队, 不回 unsupported - gcc 隔离单测 8 组全过 (tests/test_event_report.c), 单包 6 条实测 317B
This commit is contained in:
@@ -51,5 +51,6 @@ void iot_mqtt_init(void); // 初始化 MQTT 连接 (TCP→CONNECT
|
||||
void iot_mqtt_poll(void); // 主循环轮询 (状态机 + 心跳 + 重连)
|
||||
void iot_mqtt_handle_sock_int(uint8_t socketid, uint8_t intstat); // Socket 中断处理
|
||||
void iot_mqtt_publish_sensor(void); // 推送传感器数据到 MQTT
|
||||
void iot_evt_handle_ack(uint32_t msg_id, int code); // event_report 平台应答入口 (V1.04)
|
||||
|
||||
#endif /* __IOT_MQTT_SRV_H__ */
|
||||
|
||||
@@ -53,6 +53,196 @@ 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;
|
||||
}
|
||||
}
|
||||
|
||||
/* lup 传感回调: uart_srv 消费路径的帧从这里喂入 (已过校验) */
|
||||
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_evt_feed(&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_pend_id = 0;
|
||||
_evt_gaveup = 0;
|
||||
}
|
||||
_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 = now / 1000;
|
||||
_evt_retry = 0;
|
||||
iot_evt_send(_evt_pend_id, _evt_pend_ts, _evt_pend_n);
|
||||
_evt_sent_ms = now;
|
||||
}
|
||||
|
||||
/*===========================================================================
|
||||
* Buffers
|
||||
*===========================================================================*/
|
||||
@@ -382,10 +572,10 @@ static void iot_process_recv(void) {
|
||||
*===========================================================================*/
|
||||
void iot_mqtt_publish_sensor(void) {
|
||||
static uint8_t _last_car_state[4] = {0};
|
||||
if (g_iot_state != IOT_STATE_READY) return;
|
||||
if (!g_report_cfg.enable) return;
|
||||
|
||||
/*--- Step 1: 消费 0xC0 帧 → 更新缓存 ---*/
|
||||
/*--- Step 1: 消费 0xC0 帧 → 事件沿检测 + 更新缓存 ---*/
|
||||
/* 注意: 置于 READY/enable 检查之前 —— 断网/未使能期间事件照样入队,
|
||||
重连后由 iot_evt_process() 补报 (协议 V1.04) */
|
||||
if (g_pkg_uart_2.flag != 0
|
||||
&& g_pkg_uart_2.pkg[0] == 0x7F
|
||||
&& g_pkg_uart_2.pkg[3] == 0xC0) {
|
||||
@@ -396,6 +586,7 @@ void iot_mqtt_publish_sensor(void) {
|
||||
InitPkgUart(&g_pkg_uart_2); // 读完即清, 不持有
|
||||
|
||||
if (ret == 0) {
|
||||
iot_evt_feed(&sr); // 事件沿检测 (Step1 消费路径)
|
||||
memcpy(&_cached_sr, &sr, sizeof(sr));
|
||||
_cached_sr_valid = 1;
|
||||
} else {
|
||||
@@ -403,6 +594,12 @@ void iot_mqtt_publish_sensor(void) {
|
||||
}
|
||||
}
|
||||
|
||||
/*--- 事件上报状态机: 发送/重发/重连补报 ---*/
|
||||
iot_evt_process();
|
||||
|
||||
if (g_iot_state != IOT_STATE_READY) return;
|
||||
if (!g_report_cfg.enable) return;
|
||||
|
||||
/* 尚无缓存数据 → 不发 */
|
||||
if (!_cached_sr_valid) return;
|
||||
|
||||
@@ -416,7 +613,7 @@ void iot_mqtt_publish_sensor(void) {
|
||||
// variation 为有符号(V1.05): 正=车/裕量, 负=反向漂移, 两向显著变化都加速
|
||||
int32_t v = _cached_sr.coils[i].variation;
|
||||
int32_t av = (v >= 0) ? v : -v;
|
||||
if (av >= IOT_MQTT_VARIATION_THRESHOLD) {
|
||||
if (av > IOT_MQTT_VARIATION_THRESHOLD) {
|
||||
fast_mode = 1;
|
||||
break;
|
||||
}
|
||||
@@ -666,8 +863,8 @@ void iot_mqtt_init(void) {
|
||||
|
||||
PRINT("IOT: dev_serial=%s\n", g_iot_dev_serial);
|
||||
|
||||
/* 初始化传感器回调 — 用于 MQTT 上报 */
|
||||
lup_set_sensor_callback(NULL); // MQTT 模式不需要 JSON sensor_cb
|
||||
/* 初始化传感器回调 — uart_srv 消费路径的 0xC0 帧喂入事件沿检测 (V1.04) */
|
||||
lup_set_sensor_callback(iot_evt_sensor_cb);
|
||||
|
||||
_iot_reconnect_backoff = 0;
|
||||
_iot_reconnect_deadline = mstick() + 2000; // 启动后 2s 开始首次连接
|
||||
|
||||
@@ -1522,6 +1522,16 @@ void manage_mqtt_recv_message(char * msg, int length)
|
||||
return;
|
||||
}
|
||||
|
||||
// --- event_report 平台应答 (V1.04): 回显帧, 确认后出队, 不再回复 ---
|
||||
if(strcmp(cmd_str, "event_report") == 0) {
|
||||
int code = -1;
|
||||
memset(tmp, 0, sizeof(tmp));
|
||||
simple_parse_json(msg, "\"code\"", tmp);
|
||||
if(strlen(tmp) > 0) code = (int)strtol(tmp, NULL, 10);
|
||||
iot_evt_handle_ack(msg_id, code);
|
||||
return; // 应答帧终止于此, 严禁再回 unsupported (否则与平台互相打乒乓)
|
||||
}
|
||||
|
||||
// --- 未支持的命令 ---
|
||||
{
|
||||
char resp[256];
|
||||
|
||||
Reference in New Issue
Block a user