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)) }