mqtt.go 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596
  1. package initialize
  2. import (
  3. "time"
  4. mqtt "github.com/eclipse/paho.mqtt.golang"
  5. "go.uber.org/zap"
  6. "wails-app/internal/global"
  7. "wails-app/internal/service/devicebus"
  8. "wails-app/internal/service/mqttbroker"
  9. parkingService "wails-app/internal/service/parking"
  10. "wails-app/internal/service/uhf"
  11. )
  12. // InitDeviceBus 在启用 MQTT 时(可选)启动内嵌 broker、连接并装配道闸路由控制器与设备总线。
  13. // 默认关闭(mqtt.enabled=false),不影响现有 UHF 继电器/模拟道闸链路。
  14. func InitDeviceBus() {
  15. cfg := global.GVA_CONFIG.Mqtt
  16. if !cfg.Enabled {
  17. return
  18. }
  19. // 模拟道闸模式下禁止建立真实 MQTT 控制链路,避免测试指令误发到现场设备。
  20. if global.GVA_CONFIG.System.GateSimulator {
  21. if global.GVA_LOG != nil {
  22. global.GVA_LOG.Warn("道闸模拟模式已启用,跳过 MQTT 设备总线初始化")
  23. }
  24. return
  25. }
  26. // 内嵌 broker:随应用进程启动,设备直连本地端口,无需外部 mosquitto。
  27. if cfg.EmbedBroker {
  28. if _, err := mqttbroker.Start(cfg.BrokerListen); err != nil {
  29. global.GVA_LOG.Error("启动内嵌 MQTT broker 失败", zap.Error(err))
  30. return
  31. }
  32. global.GVA_LOG.Info("内嵌 MQTT broker 已启动", zap.String("listen", cfg.BrokerListen))
  33. }
  34. brokerAddr := cfg.Broker
  35. if brokerAddr == "" {
  36. brokerAddr = "tcp://127.0.0.1:1883"
  37. }
  38. clientID := cfg.ClientID
  39. if clientID == "" {
  40. clientID = "smart-parking"
  41. }
  42. opts := mqtt.NewClientOptions().
  43. AddBroker(brokerAddr).
  44. SetClientID(clientID).
  45. SetAutoReconnect(true).
  46. SetConnectRetry(true).
  47. SetConnectRetryInterval(3*time.Second).
  48. SetMaxReconnectInterval(15*time.Second).
  49. SetWill("parking/device/lwt", `{"schema":"gate.lwt.v1","device_code":"smart-parking"}`, 1, true)
  50. if cfg.Username != "" {
  51. opts.SetUsername(cfg.Username)
  52. opts.SetPassword(cfg.Password)
  53. }
  54. // OnConnect 会在首次连接和自动重连后执行,因此订阅不会因 CleanSession
  55. // 或 Broker 重启而丢失。
  56. mqttGate := parkingService.NewMQTTGateController(nil)
  57. bus := devicebus.NewDeviceBus()
  58. opts.SetOnConnectHandler(func(connected mqtt.Client) {
  59. if err := mqttGate.LoadRoutes(); err != nil {
  60. global.GVA_LOG.Error("加载 MQTT 道闸路由失败", zap.Error(err))
  61. return
  62. }
  63. if err := mqttGate.Subscribe(); err != nil {
  64. global.GVA_LOG.Error("订阅 MQTT 道闸主题失败", zap.Error(err))
  65. return
  66. }
  67. if err := bus.Subscribe(connected); err != nil {
  68. global.GVA_LOG.Error("订阅设备总线主题失败", zap.Error(err))
  69. return
  70. }
  71. global.GVA_LOG.Info("MQTT 连接已恢复,设备主题订阅完成", zap.String("broker", brokerAddr))
  72. })
  73. opts.SetConnectionLostHandler(func(_ mqtt.Client, err error) {
  74. global.GVA_LOG.Warn("MQTT 连接中断,等待自动重连", zap.Error(err))
  75. })
  76. client := mqtt.NewClient(opts)
  77. mqttGate.SetClient(client)
  78. if err := mqttGate.MarkConfiguredDevicesOffline(); err != nil {
  79. global.GVA_LOG.Error("重置 MQTT 设备状态失败", zap.Error(err))
  80. }
  81. // 初次连接不可用时保持后台重试,避免一次启动失败后永久失去设备状态同步。
  82. client.Connect()
  83. // 组合默认后端(模拟/UHF 继电器)与 MQTT 网络道闸,按设备路由。
  84. parkingService.SetGateController(parkingService.NewRoutingGateController(uhf.DeviceManager, mqttGate))
  85. global.GVA_LOG.Info("MQTT 设备总线已启用", zap.String("broker", brokerAddr))
  86. }