mqtt_integration_test.go 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231
  1. package devicebus
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "net"
  6. "testing"
  7. "time"
  8. mqtt "github.com/eclipse/paho.mqtt.golang"
  9. "github.com/glebarez/sqlite"
  10. "github.com/stretchr/testify/require"
  11. "go.uber.org/zap"
  12. "gorm.io/gorm"
  13. "wails-app/internal/dao"
  14. "wails-app/internal/global"
  15. incidentService "wails-app/internal/modules/incident/service"
  16. "wails-app/internal/service/mqttbroker"
  17. parkingService "wails-app/internal/service/parking"
  18. )
  19. func mqttE2EClient(t *testing.T, addr, id string) mqtt.Client {
  20. t.Helper()
  21. client := mqtt.NewClient(mqtt.NewClientOptions().AddBroker(addr).SetClientID(id).SetConnectTimeout(3 * time.Second))
  22. token := client.Connect()
  23. require.True(t, token.Wait())
  24. require.NoError(t, token.Error())
  25. t.Cleanup(func() { client.Disconnect(100) })
  26. return client
  27. }
  28. func mqttE2EPublish(t *testing.T, client mqtt.Client, topic string, retained bool, payload map[string]interface{}) {
  29. t.Helper()
  30. data, err := json.Marshal(payload)
  31. require.NoError(t, err)
  32. token := client.Publish(topic, 1, retained, data)
  33. require.True(t, token.Wait())
  34. require.NoError(t, token.Error())
  35. }
  36. func mqttE2EWait(t *testing.T, check func() bool) {
  37. t.Helper()
  38. require.Eventually(t, check, 8*time.Second, 25*time.Millisecond)
  39. }
  40. // TestMQTTFullPassageLifecycle 验证道闸状态和车牌识别消息经过 MQTT 后的完整业务链路。
  41. func TestMQTTFullPassageLifecycle(t *testing.T) {
  42. const (
  43. plate = "MQTT-E2E-01"
  44. inDeviceCode = "MQTT-IN-01"
  45. outDeviceCode = "MQTT-OUT-01"
  46. )
  47. dsn := fmt.Sprintf("file:mqtt-full-%d?mode=memory&cache=shared", time.Now().UnixNano())
  48. db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{})
  49. require.NoError(t, err)
  50. sqlDB, err := db.DB()
  51. require.NoError(t, err)
  52. sqlDB.SetMaxOpenConns(5)
  53. global.GVA_DB = db
  54. global.GVA_LOG = zap.NewNop()
  55. t.Cleanup(func() {
  56. parkingService.SetGateController(nil)
  57. global.GVA_DB = nil
  58. global.GVA_LOG = nil
  59. _ = sqlDB.Close()
  60. })
  61. require.NoError(t, db.AutoMigrate(
  62. &dao.ParkingLot{}, &dao.Booth{}, &dao.Channel{}, &dao.UHFReader{},
  63. &dao.VehicleType{}, &dao.Vehicle{}, &dao.VehicleRecord{}, &dao.DigitalTicket{},
  64. &dao.FeeConfig{}, &dao.PaymentRecord{}, &dao.PaymentEntryConfig{}, &dao.MonthlyCard{},
  65. &dao.Shortlist{}, &dao.IncidentRecord{}, &dao.DeviceCommandLog{},
  66. ))
  67. lot := dao.ParkingLot{LotCode: "MQTT-E2E-LOT", LotName: "MQTT全链路停车场", Capacity: 10, Available: 10}
  68. require.NoError(t, db.Create(&lot).Error)
  69. booth := dao.Booth{BoothCode: "MQTT-E2E-BOOTH", BoothName: "MQTT测试岗亭", ParkingLotID: lot.ID}
  70. require.NoError(t, db.Create(&booth).Error)
  71. inChannel := dao.Channel{ChannelCode: "MQTT-E2E-IN", ChannelName: "MQTT入口通道", Direction: "in", AllowTemporary: true, ParkingLotID: lot.ID, BoothID: booth.ID}
  72. outChannel := dao.Channel{ChannelCode: "MQTT-E2E-OUT", ChannelName: "MQTT出口通道", Direction: "out", AllowTemporary: true, ParkingLotID: lot.ID, BoothID: booth.ID}
  73. require.NoError(t, db.Create(&inChannel).Error)
  74. require.NoError(t, db.Create(&outChannel).Error)
  75. 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)
  76. 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)
  77. vehicleType := dao.VehicleType{Name: "普通车", IsSystem: false}
  78. require.NoError(t, db.Create(&vehicleType).Error)
  79. 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)
  80. require.NoError(t, db.Create(&dao.Vehicle{PlateNumber: plate, VehicleTypeID: vehicleType.ID}).Error)
  81. require.NoError(t, db.Create(&dao.PaymentEntryConfig{Code: "counter", Name: "收费岗亭", Enabled: true, IsSystem: true}).Error)
  82. // 先申请一个空闲 TCP 端口,再启动内嵌 Broker,避免和开发环境端口冲突。
  83. listener, err := net.Listen("tcp", "127.0.0.1:0")
  84. require.NoError(t, err)
  85. listenAddr := listener.Addr().String()
  86. require.NoError(t, listener.Close())
  87. broker, err := mqttbroker.Start(listenAddr)
  88. require.NoError(t, err)
  89. t.Cleanup(func() { _ = broker.Close() })
  90. brokerAddr := "tcp://" + listenAddr
  91. systemClient := mqttE2EClient(t, brokerAddr, "mqtt-e2e-system")
  92. deviceClient := mqttE2EClient(t, brokerAddr, "mqtt-e2e-devices")
  93. observerClient := mqttE2EClient(t, brokerAddr, "mqtt-e2e-observer")
  94. topic := func(channelID uint, suffix string) string {
  95. return fmt.Sprintf("parking/lot/%d/booth/%d/channel/%d/%s", lot.ID, booth.ID, channelID, suffix)
  96. }
  97. mqttGate := parkingService.NewMQTTGateController(systemClient)
  98. require.NoError(t, mqttGate.LoadRoutes())
  99. require.NoError(t, mqttGate.Subscribe())
  100. parkingService.SetGateController(mqttGate)
  101. bus := NewDeviceBus()
  102. require.NoError(t, bus.Subscribe(systemClient))
  103. cameraEvents := make(chan string, 4)
  104. observerToken := observerClient.Subscribe("parking/lot/+/booth/+/channel/+/camera/event", 1, func(_ mqtt.Client, msg mqtt.Message) {
  105. cameraEvents <- msg.Topic()
  106. })
  107. require.True(t, observerToken.Wait())
  108. require.NoError(t, observerToken.Error())
  109. // 模拟设备:收到系统命令后,向同一通道发布成功 ACK。
  110. ackReceived := make(chan string, 4)
  111. ack := func(_ mqtt.Client, msg mqtt.Message) {
  112. var command struct {
  113. CmdID string `json:"cmd_id"`
  114. DeviceCode string `json:"device_code"`
  115. }
  116. if json.Unmarshal(msg.Payload(), &command) != nil || command.CmdID == "" {
  117. return
  118. }
  119. deviceCode := command.DeviceCode
  120. select {
  121. case ackReceived <- deviceCode:
  122. default:
  123. }
  124. channelID := inChannel.ID
  125. if deviceCode == outDeviceCode {
  126. channelID = outChannel.ID
  127. }
  128. data, _ := json.Marshal(map[string]interface{}{"schema": "gate.ack.v1", "cmd_id": command.CmdID, "device_code": deviceCode, "result": "success", "error": ""})
  129. token := deviceClient.Publish(topic(channelID, "gate/cmd/ack"), 1, false, data)
  130. token.Wait()
  131. }
  132. subscribeToken := deviceClient.SubscribeMultiple(map[string]byte{
  133. topic(inChannel.ID, "gate/cmd"): 1,
  134. topic(outChannel.ID, "gate/cmd"): 1,
  135. }, ack)
  136. require.True(t, subscribeToken.Wait())
  137. require.NoError(t, subscribeToken.Error())
  138. // 预检 MQTT 命令和 ACK 通路,避免把协议失败误判为车辆业务失败。
  139. openDone := make(chan error, 1)
  140. go func() { openDone <- mqttGate.OpenGate(inDeviceCode, 1) }()
  141. mqttE2EWait(t, func() bool {
  142. select {
  143. case deviceCode := <-ackReceived:
  144. return deviceCode == inDeviceCode
  145. default:
  146. return false
  147. }
  148. })
  149. require.NoError(t, <-openDone)
  150. // 1. 两个道闸上线,验证内存状态和设备表状态均更新。
  151. mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "gate/state"), true, map[string]interface{}{"schema": "gate.state.v1", "device_code": inDeviceCode, "state": "closed"})
  152. mqttE2EPublish(t, deviceClient, topic(outChannel.ID, "gate/state"), true, map[string]interface{}{"schema": "gate.state.v1", "device_code": outDeviceCode, "state": "closed"})
  153. mqttE2EWait(t, func() bool {
  154. var count int64
  155. return db.Model(&dao.UHFReader{}).Where("status = ?", "online").Count(&count).Error == nil && count == 2
  156. })
  157. require.True(t, mqttGate.IsGateConnected(inDeviceCode))
  158. require.True(t, mqttGate.IsGateConnected(outDeviceCode))
  159. // 2. 入口道闸离线后恢复上线,同时确认产生离线异常记录。
  160. mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "gate/lwt"), true, map[string]interface{}{"schema": "gate.lwt.v1", "device_code": inDeviceCode})
  161. mqttE2EWait(t, func() bool {
  162. var device dao.UHFReader
  163. return db.Where("device_code = ?", inDeviceCode).First(&device).Error == nil && device.Status == "offline"
  164. })
  165. require.False(t, mqttGate.IsGateConnected(inDeviceCode))
  166. var offlineIncidents int64
  167. mqttE2EWait(t, func() bool {
  168. return db.Model(&dao.IncidentRecord{}).Where("category = ? AND device_code = ?", incidentService.CategoryDeviceOffline, inDeviceCode).Count(&offlineIncidents).Error == nil && offlineIncidents >= 1
  169. })
  170. mqttE2EPublish(t, deviceClient, topic(inChannel.ID, "gate/state"), true, map[string]interface{}{"schema": "gate.state.v1", "device_code": inDeviceCode, "state": "closed"})
  171. mqttE2EWait(t, func() bool {
  172. var device dao.UHFReader
  173. return db.Where("device_code = ?", inDeviceCode).First(&device).Error == nil && device.Status == "online"
  174. })
  175. // 3. 车牌识别入场,验证停车会话、数字票、余位和入口开闸流水。
  176. 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()})
  177. mqttE2EWait(t, func() bool { return len(cameraEvents) >= 1 })
  178. var record dao.VehicleRecord
  179. mqttE2EWait(t, func() bool {
  180. return db.Where("plate_number = ? AND exit_time IS NULL", plate).First(&record).Error == nil
  181. })
  182. var ticket dao.DigitalTicket
  183. require.NoError(t, db.Where("vehicle_record_id = ?", record.ID).First(&ticket).Error)
  184. require.Equal(t, "pending_payment", ticket.State)
  185. require.Equal(t, inChannel.ID, record.EntryChannelID)
  186. require.Equal(t, inDeviceCode, record.EntryDeviceCode)
  187. require.Equal(t, 9, lotCapacity(t, db, lot.ID))
  188. var entryCommand dao.DeviceCommandLog
  189. require.Eventually(t, func() bool {
  190. return db.Where("device_code = ?", inDeviceCode).Order("id DESC").First(&entryCommand).Error == nil && entryCommand.Result == incidentService.CommandSuccess
  191. }, 8*time.Second, 25*time.Millisecond)
  192. require.Equal(t, record.ID, entryCommand.SessionID)
  193. // 4. 车牌识别出场,验证支付状态、会话关闭、余位恢复和出口开闸流水关联。
  194. // 统一通行入口按秒记录三秒防抖;等待四秒确保跨过时间边界。
  195. time.Sleep(4 * time.Second)
  196. 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()})
  197. mqttE2EWait(t, func() bool { return len(cameraEvents) >= 2 })
  198. mqttE2EWait(t, func() bool { return db.First(&record, record.ID).Error == nil && record.ExitTime != nil })
  199. require.Equal(t, "paid", record.PaymentStatus)
  200. require.Equal(t, outChannel.ID, record.ExitChannelID)
  201. require.Equal(t, outDeviceCode, record.ExitDeviceCode)
  202. require.NoError(t, db.First(&ticket, ticket.ID).Error)
  203. require.Equal(t, "exited", ticket.State)
  204. require.Equal(t, 10, lotCapacity(t, db, lot.ID))
  205. var exitCommand dao.DeviceCommandLog
  206. require.NoError(t, db.Where("device_code = ?", outDeviceCode).Order("id DESC").First(&exitCommand).Error)
  207. require.Equal(t, incidentService.CommandSuccess, exitCommand.Result)
  208. require.Equal(t, record.ID, exitCommand.SessionID)
  209. }
  210. func lotCapacity(t *testing.T, db *gorm.DB, id uint) int {
  211. t.Helper()
  212. var lot dao.ParkingLot
  213. require.NoError(t, db.First(&lot, id).Error)
  214. return lot.Available
  215. }