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