| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596 |
- package initialize
- import (
- "time"
- mqtt "github.com/eclipse/paho.mqtt.golang"
- "go.uber.org/zap"
- "wails-app/internal/global"
- "wails-app/internal/service/devicebus"
- "wails-app/internal/service/mqttbroker"
- parkingService "wails-app/internal/service/parking"
- "wails-app/internal/service/uhf"
- )
- // InitDeviceBus 在启用 MQTT 时(可选)启动内嵌 broker、连接并装配道闸路由控制器与设备总线。
- // 默认关闭(mqtt.enabled=false),不影响现有 UHF 继电器/模拟道闸链路。
- func InitDeviceBus() {
- cfg := global.GVA_CONFIG.Mqtt
- if !cfg.Enabled {
- return
- }
- // 模拟道闸模式下禁止建立真实 MQTT 控制链路,避免测试指令误发到现场设备。
- if global.GVA_CONFIG.System.GateSimulator {
- if global.GVA_LOG != nil {
- global.GVA_LOG.Warn("道闸模拟模式已启用,跳过 MQTT 设备总线初始化")
- }
- return
- }
- // 内嵌 broker:随应用进程启动,设备直连本地端口,无需外部 mosquitto。
- if cfg.EmbedBroker {
- if _, err := mqttbroker.Start(cfg.BrokerListen); err != nil {
- global.GVA_LOG.Error("启动内嵌 MQTT broker 失败", zap.Error(err))
- return
- }
- global.GVA_LOG.Info("内嵌 MQTT broker 已启动", zap.String("listen", cfg.BrokerListen))
- }
- brokerAddr := cfg.Broker
- if brokerAddr == "" {
- brokerAddr = "tcp://127.0.0.1:1883"
- }
- clientID := cfg.ClientID
- if clientID == "" {
- clientID = "smart-parking"
- }
- opts := mqtt.NewClientOptions().
- AddBroker(brokerAddr).
- SetClientID(clientID).
- SetAutoReconnect(true).
- SetConnectRetry(true).
- SetConnectRetryInterval(3*time.Second).
- SetMaxReconnectInterval(15*time.Second).
- SetWill("parking/device/lwt", `{"schema":"gate.lwt.v1","device_code":"smart-parking"}`, 1, true)
- if cfg.Username != "" {
- opts.SetUsername(cfg.Username)
- opts.SetPassword(cfg.Password)
- }
- // OnConnect 会在首次连接和自动重连后执行,因此订阅不会因 CleanSession
- // 或 Broker 重启而丢失。
- mqttGate := parkingService.NewMQTTGateController(nil)
- bus := devicebus.NewDeviceBus()
- opts.SetOnConnectHandler(func(connected mqtt.Client) {
- if err := mqttGate.LoadRoutes(); err != nil {
- global.GVA_LOG.Error("加载 MQTT 道闸路由失败", zap.Error(err))
- return
- }
- if err := mqttGate.Subscribe(); err != nil {
- global.GVA_LOG.Error("订阅 MQTT 道闸主题失败", zap.Error(err))
- return
- }
- if err := bus.Subscribe(connected); err != nil {
- global.GVA_LOG.Error("订阅设备总线主题失败", zap.Error(err))
- return
- }
- global.GVA_LOG.Info("MQTT 连接已恢复,设备主题订阅完成", zap.String("broker", brokerAddr))
- })
- opts.SetConnectionLostHandler(func(_ mqtt.Client, err error) {
- global.GVA_LOG.Warn("MQTT 连接中断,等待自动重连", zap.Error(err))
- })
- client := mqtt.NewClient(opts)
- mqttGate.SetClient(client)
- if err := mqttGate.MarkConfiguredDevicesOffline(); err != nil {
- global.GVA_LOG.Error("重置 MQTT 设备状态失败", zap.Error(err))
- }
- // 初次连接不可用时保持后台重试,避免一次启动失败后永久失去设备状态同步。
- client.Connect()
- // 组合默认后端(模拟/UHF 继电器)与 MQTT 网络道闸,按设备路由。
- parkingService.SetGateController(parkingService.NewRoutingGateController(uhf.DeviceManager, mqttGate))
- global.GVA_LOG.Info("MQTT 设备总线已启用", zap.String("broker", brokerAddr))
- }
|