package devicebus import ( "encoding/json" "fmt" "net" "testing" "time" mqtt "github.com/eclipse/paho.mqtt.golang" "github.com/glebarez/sqlite" "github.com/stretchr/testify/require" "go.uber.org/zap" "gorm.io/gorm" "wails-app/internal/dao" "wails-app/internal/global" incidentService "wails-app/internal/modules/incident/service" "wails-app/internal/service/mqttbroker" parkingService "wails-app/internal/service/parking" ) func mqttE2EClient(t *testing.T, addr, id string) mqtt.Client { t.Helper() client := mqtt.NewClient(mqtt.NewClientOptions().AddBroker(addr).SetClientID(id).SetConnectTimeout(3 * time.Second)) token := client.Connect() require.True(t, token.Wait()) require.NoError(t, token.Error()) t.Cleanup(func() { client.Disconnect(100) }) return client } func mqttE2EPublish(t *testing.T, client mqtt.Client, topic string, retained bool, payload map[string]interface{}) { t.Helper() data, err := json.Marshal(payload) require.NoError(t, err) token := client.Publish(topic, 1, retained, data) require.True(t, token.Wait()) require.NoError(t, token.Error()) } func mqttE2EWait(t *testing.T, check func() bool) { t.Helper() require.Eventually(t, check, 8*time.Second, 25*time.Millisecond) } // TestMQTTFullPassageLifecycle 验证道闸状态和车牌识别消息经过 MQTT 后的完整业务链路。 func TestMQTTFullPassageLifecycle(t *testing.T) { const ( plate = "MQTT-E2E-01" inDeviceCode = "MQTT-IN-01" outDeviceCode = "MQTT-OUT-01" ) dsn := fmt.Sprintf("file:mqtt-full-%d?mode=memory&cache=shared", time.Now().UnixNano()) db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{}) require.NoError(t, err) sqlDB, err := db.DB() require.NoError(t, err) sqlDB.SetMaxOpenConns(5) global.GVA_DB = db global.GVA_LOG = zap.NewNop() t.Cleanup(func() { parkingService.SetGateController(nil) global.GVA_DB = nil global.GVA_LOG = nil _ = sqlDB.Close() }) require.NoError(t, db.AutoMigrate( &dao.ParkingLot{}, &dao.Booth{}, &dao.Channel{}, &dao.UHFReader{}, &dao.VehicleType{}, &dao.Vehicle{}, &dao.VehicleRecord{}, &dao.DigitalTicket{}, &dao.FeeConfig{}, &dao.PaymentRecord{}, &dao.PaymentEntryConfig{}, &dao.MonthlyCard{}, &dao.Shortlist{}, &dao.IncidentRecord{}, &dao.DeviceCommandLog{}, )) lot := dao.ParkingLot{LotCode: "MQTT-E2E-LOT", LotName: "MQTT全链路停车场", Capacity: 10, Available: 10} require.NoError(t, db.Create(&lot).Error) booth := dao.Booth{BoothCode: "MQTT-E2E-BOOTH", BoothName: "MQTT测试岗亭", ParkingLotID: lot.ID} require.NoError(t, db.Create(&booth).Error) inChannel := dao.Channel{ChannelCode: "MQTT-E2E-IN", ChannelName: "MQTT入口通道", Direction: "in", AllowTemporary: true, ParkingLotID: lot.ID, BoothID: booth.ID} outChannel := dao.Channel{ChannelCode: "MQTT-E2E-OUT", ChannelName: "MQTT出口通道", Direction: "out", AllowTemporary: true, ParkingLotID: lot.ID, BoothID: booth.ID} require.NoError(t, db.Create(&inChannel).Error) require.NoError(t, db.Create(&outChannel).Error) require.NoError(t, db.Create(&dao.UHFReader{DeviceCode: inDeviceCode, DeviceName: "MQTT入口道闸", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, IsActive: true, ParkingLotID: lot.ID, ChannelID: inChannel.ID}).Error) require.NoError(t, db.Create(&dao.UHFReader{DeviceCode: outDeviceCode, DeviceName: "MQTT出口道闸", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, IsActive: true, ParkingLotID: lot.ID, ChannelID: outChannel.ID}).Error) vehicleType := dao.VehicleType{Name: "普通车", IsSystem: false} require.NoError(t, db.Create(&vehicleType).Error) require.NoError(t, db.Create(&dao.FeeConfig{VehicleTypeID: vehicleType.ID, UnitTime: 60, StartTime: 0, StartFee: 0, UnitFee: 0, DailyMaxFee: 0, FreeExitMinutes: 15, VIPDiscount: 1}).Error) require.NoError(t, db.Create(&dao.Vehicle{PlateNumber: plate, VehicleTypeID: vehicleType.ID}).Error) require.NoError(t, db.Create(&dao.PaymentEntryConfig{Code: "counter", Name: "收费岗亭", Enabled: true, IsSystem: true}).Error) // 先申请一个空闲 TCP 端口,再启动内嵌 Broker,避免和开发环境端口冲突。 listener, err := net.Listen("tcp", "127.0.0.1:0") require.NoError(t, err) listenAddr := listener.Addr().String() require.NoError(t, listener.Close()) broker, err := mqttbroker.Start(listenAddr) require.NoError(t, err) t.Cleanup(func() { _ = broker.Close() }) brokerAddr := "tcp://" + listenAddr systemClient := mqttE2EClient(t, brokerAddr, "mqtt-e2e-system") deviceClient := mqttE2EClient(t, brokerAddr, "mqtt-e2e-devices") observerClient := mqttE2EClient(t, brokerAddr, "mqtt-e2e-observer") topic := func(channelID uint, suffix string) string { return fmt.Sprintf("parking/lot/%d/booth/%d/channel/%d/%s", lot.ID, booth.ID, channelID, suffix) } mqttGate := parkingService.NewMQTTGateController(systemClient) require.NoError(t, mqttGate.LoadRoutes()) require.NoError(t, mqttGate.Subscribe()) parkingService.SetGateController(mqttGate) bus := NewDeviceBus() require.NoError(t, bus.Subscribe(systemClient)) cameraEvents := make(chan string, 4) observerToken := observerClient.Subscribe("parking/lot/+/booth/+/channel/+/camera/event", 1, func(_ mqtt.Client, msg mqtt.Message) { cameraEvents <- msg.Topic() }) require.True(t, observerToken.Wait()) require.NoError(t, observerToken.Error()) // 模拟设备:收到系统命令后,向同一通道发布成功 ACK。 ackReceived := make(chan string, 4) ack := func(_ mqtt.Client, msg mqtt.Message) { var command struct { CmdID string `json:"cmd_id"` DeviceCode string `json:"device_code"` } if json.Unmarshal(msg.Payload(), &command) != nil || command.CmdID == "" { return } deviceCode := command.DeviceCode select { case ackReceived <- deviceCode: default: } channelID := inChannel.ID if deviceCode == outDeviceCode { channelID = outChannel.ID } data, _ := json.Marshal(map[string]interface{}{"schema": "gate.ack.v1", "cmd_id": command.CmdID, "device_code": deviceCode, "result": "success", "error": ""}) token := deviceClient.Publish(topic(channelID, "gate/cmd/ack"), 1, false, data) token.Wait() } subscribeToken := deviceClient.SubscribeMultiple(map[string]byte{ topic(inChannel.ID, "gate/cmd"): 1, topic(outChannel.ID, "gate/cmd"): 1, }, ack) require.True(t, subscribeToken.Wait()) require.NoError(t, subscribeToken.Error()) // 预检 MQTT 命令和 ACK 通路,避免把协议失败误判为车辆业务失败。 openDone := make(chan error, 1) go func() { openDone <- mqttGate.OpenGate(inDeviceCode, 1) }() mqttE2EWait(t, func() bool { select { case deviceCode := <-ackReceived: return deviceCode == inDeviceCode default: return false } }) require.NoError(t, <-openDone) // 1. 两个道闸上线,验证内存状态和设备表状态均更新。 mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "gate/state"), true, map[string]interface{}{"schema": "gate.state.v1", "device_code": inDeviceCode, "state": "closed"}) mqttE2EPublish(t, deviceClient, topic(outChannel.ID, "gate/state"), true, map[string]interface{}{"schema": "gate.state.v1", "device_code": outDeviceCode, "state": "closed"}) mqttE2EWait(t, func() bool { var count int64 return db.Model(&dao.UHFReader{}).Where("status = ?", "online").Count(&count).Error == nil && count == 2 }) require.True(t, mqttGate.IsGateConnected(inDeviceCode)) require.True(t, mqttGate.IsGateConnected(outDeviceCode)) // 2. 入口道闸离线后恢复上线,同时确认产生离线异常记录。 mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "gate/lwt"), true, map[string]interface{}{"schema": "gate.lwt.v1", "device_code": inDeviceCode}) mqttE2EWait(t, func() bool { var device dao.UHFReader return db.Where("device_code = ?", inDeviceCode).First(&device).Error == nil && device.Status == "offline" }) require.False(t, mqttGate.IsGateConnected(inDeviceCode)) var offlineIncidents int64 mqttE2EWait(t, func() bool { return db.Model(&dao.IncidentRecord{}).Where("category = ? AND device_code = ?", incidentService.CategoryDeviceOffline, inDeviceCode).Count(&offlineIncidents).Error == nil && offlineIncidents >= 1 }) mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "gate/state"), true, map[string]interface{}{"schema": "gate.state.v1", "device_code": inDeviceCode, "state": "closed"}) mqttE2EWait(t, func() bool { var device dao.UHFReader return db.Where("device_code = ?", inDeviceCode).First(&device).Error == nil && device.Status == "online" }) // 3. 车牌识别入场,验证停车会话、数字票、余位和入口开闸流水。 mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "camera/event"), false, map[string]interface{}{"schema": "camera.event.v1", "device_code": inDeviceCode, "plate_number": plate, "direction": "in", "ts": time.Now().Unix()}) mqttE2EWait(t, func() bool { return len(cameraEvents) >= 1 }) var record dao.VehicleRecord mqttE2EWait(t, func() bool { return db.Where("plate_number = ? AND exit_time IS NULL", plate).First(&record).Error == nil }) var ticket dao.DigitalTicket require.NoError(t, db.Where("vehicle_record_id = ?", record.ID).First(&ticket).Error) require.Equal(t, "pending_payment", ticket.State) require.Equal(t, inChannel.ID, record.EntryChannelID) require.Equal(t, inDeviceCode, record.EntryDeviceCode) require.Equal(t, 9, lotCapacity(t, db, lot.ID)) var entryCommand dao.DeviceCommandLog require.Eventually(t, func() bool { return db.Where("device_code = ?", inDeviceCode).Order("id DESC").First(&entryCommand).Error == nil && entryCommand.Result == incidentService.CommandSuccess }, 8*time.Second, 25*time.Millisecond) require.Equal(t, record.ID, entryCommand.SessionID) // 4. 车牌识别出场,验证支付状态、会话关闭、余位恢复和出口开闸流水关联。 // 统一通行入口按秒记录三秒防抖;等待四秒确保跨过时间边界。 time.Sleep(4 * time.Second) mqttE2EPublish(t, deviceClient, topic(outChannel.ID, "camera/event"), false, map[string]interface{}{"schema": "camera.event.v1", "device_code": outDeviceCode, "plate_number": plate, "direction": "out", "ts": time.Now().Unix()}) mqttE2EWait(t, func() bool { return len(cameraEvents) >= 2 }) mqttE2EWait(t, func() bool { return db.First(&record, record.ID).Error == nil && record.ExitTime != nil }) require.Equal(t, "paid", record.PaymentStatus) require.Equal(t, outChannel.ID, record.ExitChannelID) require.Equal(t, outDeviceCode, record.ExitDeviceCode) require.NoError(t, db.First(&ticket, ticket.ID).Error) require.Equal(t, "exited", ticket.State) require.Equal(t, 10, lotCapacity(t, db, lot.ID)) var exitCommand dao.DeviceCommandLog require.NoError(t, db.Where("device_code = ?", outDeviceCode).Order("id DESC").First(&exitCommand).Error) require.Equal(t, incidentService.CommandSuccess, exitCommand.Result) require.Equal(t, record.ID, exitCommand.SessionID) } func lotCapacity(t *testing.T, db *gorm.DB, id uint) int { t.Helper() var lot dao.ParkingLot require.NoError(t, db.First(&lot, id).Error) return lot.Available }