mqtt_business_test.py 2.8 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192
  1. """DTU MQTT 业务端模拟测试脚本
  2. 模拟上层业务系统:
  3. 1. 订阅 DTU 上行主题(register/status/response/patchpanel/jumper/env)
  4. 2. 下发全套下行控制指令到 control 主题
  5. 3. 打印 DTU 返回的响应/状态,验证协议端到端
  6. 参照《线架系统功能需求与MQTT协议设计.md》第 7 章。
  7. """
  8. import json
  9. import time
  10. import uuid
  11. import paho.mqtt.client as mqtt
  12. # ===== 配置 =====
  13. BROKER = "xt.wenhq.top"
  14. PORT = 8581
  15. USER = "admin"
  16. PASS = "admin"
  17. PREFIX = "线架系统"
  18. CUSTOMER = "default_customer"
  19. DTU_ID = "dtu_001"
  20. PANEL_ID = "PANEL_dtu_001_2" # DTU 当前在线面板(地址2)
  21. control_topic = f"{PREFIX}/{CUSTOMER}/dtu/{DTU_ID}/control"
  22. received = []
  23. def on_connect(client, userdata, flags, rc):
  24. print(f"[业务端] 已连接 broker rc={rc}")
  25. client.subscribe(f"{PREFIX}/{CUSTOMER}/#", qos=1)
  26. print(f"[业务端] 已订阅 {PREFIX}/{CUSTOMER}/#")
  27. def on_message(client, userdata, msg):
  28. try:
  29. payload = json.loads(msg.payload.decode("utf-8"))
  30. body = json.dumps(payload, ensure_ascii=False)
  31. except Exception:
  32. body = msg.payload.decode("utf-8", errors="replace")
  33. print(f"[收到] {msg.topic}\n {body}")
  34. received.append((msg.topic, payload))
  35. def send_control(client, command, target, params=None):
  36. envelope = {
  37. "msg_id": f"ctrl_{uuid.uuid4().hex[:8]}",
  38. "timestamp": int(time.time() * 1000),
  39. "dtu_id": DTU_ID,
  40. "type": "CONTROL",
  41. "payload": {"command": command, "target": target, "params": params or {}},
  42. }
  43. print(f"\n[下发] {control_topic}\n {json.dumps(envelope, ensure_ascii=False)}")
  44. client.publish(control_topic, json.dumps(envelope), qos=1)
  45. return envelope["msg_id"]
  46. def main():
  47. client = mqtt.Client(client_id="business_test")
  48. client.username_pw_set(USER, PASS)
  49. client.on_connect = on_connect
  50. client.on_message = on_message
  51. client.connect(BROKER, PORT, 60)
  52. client.loop_start()
  53. print("等待 3 秒收集上行消息(register/status)...")
  54. time.sleep(3)
  55. # 下发全套控制指令(每条间隔 3 秒等响应)
  56. send_control(client, "QUERY_DTU_STATUS", "dtu")
  57. time.sleep(3)
  58. send_control(client, "READ_PANEL_STATUS", PANEL_ID)
  59. time.sleep(3)
  60. send_control(client, "QUERY_JUMPER_STATUS", "all")
  61. time.sleep(3)
  62. send_control(client, "QUERY_ENV_SENSOR", "all")
  63. time.sleep(3)
  64. send_control(client, "SET_PORT_LED", PANEL_ID, {"port_id": 1, "led_mode": "BLINK_RED"})
  65. time.sleep(3)
  66. send_control(client, "SET_PORT_LED", PANEL_ID, {"port_id": 1, "led_mode": "OFF"})
  67. time.sleep(2)
  68. print(f"\n=== 测试完成,共收到 {len(received)} 条上行消息 ===")
  69. client.loop_stop()
  70. client.disconnect()
  71. if __name__ == "__main__":
  72. main()