refactor(vd960DBN): MQTT 重构,对标 DBN101GA 参考项目
net_srv.c 改动: - WCHNET_HandleSockInt IoT 分支:SINT_STAT_CONNECT 直接调 mqtt_connect(), SINT_STAT_RECV 调 mqtt_data_manage() — 复用已有的 SocketId_TCP/mqttBuf/mqttData - 新增 mqtt_data_manage(): 处理 CONNACK/SUBACK/PUBLISH/PINGRESP - 新增 poll_mqtt(): 心跳 PINGREQ (~9s) peripheral_main.c 改动: - IoT 模式下 poll_mqtt() 替代 iot_mqtt_poll() 对标 DBN101GA 的关键差异: - net_srv.c 已有完整 mqtt_connect/mqtt_publish/MQTT_Subscribe/MQTT_Pingreq - SocketId_TCP 由 WCHNET_CreateTcpMqttSocket 创建,不再自建 socket - SINT_STAT_CONNECT 直接发 MQTT CONNECT,不延迟
This commit is contained in:
@@ -377,6 +377,74 @@ void reset_net_active_timeup(void)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* MQTT 数据接收处理 — IoT 模式 (ref: DBN101GA) */
|
||||||
|
void mqtt_data_manage(uint8_t id)
|
||||||
|
{
|
||||||
|
uint32_t len;
|
||||||
|
static uint8_t MyBuf[RECE_BUF_LEN]; // static: 避免栈溢出
|
||||||
|
unsigned char dup;
|
||||||
|
unsigned short packetid;
|
||||||
|
int qos;
|
||||||
|
unsigned char retained;
|
||||||
|
MQTTString topicName;
|
||||||
|
unsigned char *payload;
|
||||||
|
int payloadlen;
|
||||||
|
|
||||||
|
memset(MyBuf, 0, RECE_BUF_LEN);
|
||||||
|
len = WCHNET_SocketRecvLen(id, NULL);
|
||||||
|
WCHNET_SocketRecv(id, MyBuf, &len);
|
||||||
|
PRINT("MQTT recv sock=%d len=%d type=0x%02X\n", id, (int)len, MyBuf[0]);
|
||||||
|
|
||||||
|
switch (MyBuf[0] >> 4) {
|
||||||
|
case CONNACK:
|
||||||
|
PRINT("MQTT CONNACK\n");
|
||||||
|
dg_subscribe_display_topic();
|
||||||
|
dev_initialize_pub();
|
||||||
|
reset_net_active_timeup();
|
||||||
|
break;
|
||||||
|
case PUBLISH:
|
||||||
|
MQTTDeserialize_publish(&dup, &qos, &retained, &packetid,
|
||||||
|
&topicName, &payload, &payloadlen, MyBuf, len);
|
||||||
|
manage_mqtt_recv_message((char *)payload, payloadlen);
|
||||||
|
reset_net_active_timeup();
|
||||||
|
break;
|
||||||
|
case SUBACK:
|
||||||
|
PRINT("MQTT SUBACK\n");
|
||||||
|
break;
|
||||||
|
case PINGRESP:
|
||||||
|
PRINT("MQTT PINGRESP\n");
|
||||||
|
reset_net_active_timeup();
|
||||||
|
break;
|
||||||
|
default:
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
g_net_state.flag = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
/*********************************************************************
|
||||||
|
* @fn poll_mqtt
|
||||||
|
*
|
||||||
|
* @brief MQTT 心跳轮询 — IoT 模式 (ref: DBN101GA)
|
||||||
|
*
|
||||||
|
* @return none
|
||||||
|
*/
|
||||||
|
void poll_mqtt(void)
|
||||||
|
{
|
||||||
|
if (g_net_state.flag < 2) return;
|
||||||
|
|
||||||
|
static uint8_t _ping_counter = 0;
|
||||||
|
_ping_counter++;
|
||||||
|
|
||||||
|
if (_ping_counter >= 90) { // ~9s (roughly, based on loop period)
|
||||||
|
_ping_counter = 0;
|
||||||
|
MQTT_Pingreq();
|
||||||
|
PRINT("MQTT PINGREQ\n");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
/*********************************************************************
|
/*********************************************************************
|
||||||
* @fn WCHNET_DataManage
|
* @fn WCHNET_DataManage
|
||||||
@@ -485,9 +553,27 @@ uint8_t i;
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// IoT MQTT mode -- route to MQTT handler
|
// IoT MQTT mode — inline MQTT handling (ref: DBN101GA)
|
||||||
if (g_sub_code_enable.iot_enable) {
|
if (g_sub_code_enable.iot_enable) {
|
||||||
iot_mqtt_handle_sock_int(socketid, intstat);
|
if (intstat & SINT_STAT_RECV) {
|
||||||
|
mqtt_data_manage(socketid);
|
||||||
|
}
|
||||||
|
if (intstat & SINT_STAT_CONNECT) {
|
||||||
|
#if KEEPLIVE_ENABLE
|
||||||
|
WCHNET_SocketSetKeepLive(socketid, ENABLE);
|
||||||
|
#endif
|
||||||
|
WCHNET_ModifyRecvBuf(socketid, (uint32_t)SocketRecvBuf[socketid], RECE_BUF_LEN);
|
||||||
|
PRINT("TCP Connect Success (MQTT)\n");
|
||||||
|
mqtt_connect();
|
||||||
|
}
|
||||||
|
if (intstat & SINT_STAT_DISCONNECT) {
|
||||||
|
PRINT("TCP Disconnect (MQTT)\n");
|
||||||
|
g_net_state.flag = 1;
|
||||||
|
}
|
||||||
|
if (intstat & SINT_STAT_TIM_OUT) {
|
||||||
|
PRINT("TCP Timeout (MQTT)\n");
|
||||||
|
g_net_state.flag = 1;
|
||||||
|
}
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -271,7 +271,7 @@ void Main_Circulation(void)
|
|||||||
|
|
||||||
if (g_sub_code_enable.iot_enable) {
|
if (g_sub_code_enable.iot_enable) {
|
||||||
iot_mqtt_publish_sensor(); // Push 0xC0 sensor data to MQTT broker
|
iot_mqtt_publish_sensor(); // Push 0xC0 sensor data to MQTT broker
|
||||||
iot_mqtt_poll(); // IoT MQTT state machine + heartbeat
|
poll_mqtt(); // MQTT PINGREQ heartbeat (in net_srv.c)
|
||||||
} else {
|
} else {
|
||||||
tcp_json_push_sensor(); // Push 0xC0 sensor data to TCP JSON client
|
tcp_json_push_sensor(); // Push 0xC0 sensor data to TCP JSON client
|
||||||
tcp_json_poll();
|
tcp_json_poll();
|
||||||
|
|||||||
Reference in New Issue
Block a user