| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192 |
- """DTU MQTT 业务端模拟测试脚本
- 模拟上层业务系统:
- 1. 订阅 DTU 上行主题(register/status/response/patchpanel/jumper/env)
- 2. 下发全套下行控制指令到 control 主题
- 3. 打印 DTU 返回的响应/状态,验证协议端到端
- 参照《线架系统功能需求与MQTT协议设计.md》第 7 章。
- """
- import json
- import time
- import uuid
- import paho.mqtt.client as mqtt
- # ===== 配置 =====
- BROKER = "xt.wenhq.top"
- PORT = 8581
- USER = "admin"
- PASS = "admin"
- PREFIX = "线架系统"
- CUSTOMER = "default_customer"
- DTU_ID = "dtu_001"
- PANEL_ID = "PANEL_dtu_001_2" # DTU 当前在线面板(地址2)
- control_topic = f"{PREFIX}/{CUSTOMER}/dtu/{DTU_ID}/control"
- received = []
- def on_connect(client, userdata, flags, rc):
- print(f"[业务端] 已连接 broker rc={rc}")
- client.subscribe(f"{PREFIX}/{CUSTOMER}/#", qos=1)
- print(f"[业务端] 已订阅 {PREFIX}/{CUSTOMER}/#")
- def on_message(client, userdata, msg):
- try:
- payload = json.loads(msg.payload.decode("utf-8"))
- body = json.dumps(payload, ensure_ascii=False)
- except Exception:
- body = msg.payload.decode("utf-8", errors="replace")
- print(f"[收到] {msg.topic}\n {body}")
- received.append((msg.topic, payload))
- def send_control(client, command, target, params=None):
- envelope = {
- "msg_id": f"ctrl_{uuid.uuid4().hex[:8]}",
- "timestamp": int(time.time() * 1000),
- "dtu_id": DTU_ID,
- "type": "CONTROL",
- "payload": {"command": command, "target": target, "params": params or {}},
- }
- print(f"\n[下发] {control_topic}\n {json.dumps(envelope, ensure_ascii=False)}")
- client.publish(control_topic, json.dumps(envelope), qos=1)
- return envelope["msg_id"]
- def main():
- client = mqtt.Client(client_id="business_test")
- client.username_pw_set(USER, PASS)
- client.on_connect = on_connect
- client.on_message = on_message
- client.connect(BROKER, PORT, 60)
- client.loop_start()
- print("等待 3 秒收集上行消息(register/status)...")
- time.sleep(3)
- # 下发全套控制指令(每条间隔 3 秒等响应)
- send_control(client, "QUERY_DTU_STATUS", "dtu")
- time.sleep(3)
- send_control(client, "READ_PANEL_STATUS", PANEL_ID)
- time.sleep(3)
- send_control(client, "QUERY_JUMPER_STATUS", "all")
- time.sleep(3)
- send_control(client, "QUERY_ENV_SENSOR", "all")
- time.sleep(3)
- send_control(client, "SET_PORT_LED", PANEL_ID, {"port_id": 1, "led_mode": "BLINK_RED"})
- time.sleep(3)
- send_control(client, "SET_PORT_LED", PANEL_ID, {"port_id": 1, "led_mode": "OFF"})
- time.sleep(2)
- print(f"\n=== 测试完成,共收到 {len(received)} 条上行消息 ===")
- client.loop_stop()
- client.disconnect()
- if __name__ == "__main__":
- main()
|