mqtt_gate_controller_test.go 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119
  1. package parking
  2. import (
  3. "fmt"
  4. "testing"
  5. "time"
  6. mqtt "github.com/eclipse/paho.mqtt.golang"
  7. "github.com/glebarez/sqlite"
  8. "github.com/stretchr/testify/require"
  9. "go.uber.org/zap"
  10. "gorm.io/gorm"
  11. "wails-app/internal/dao"
  12. "wails-app/internal/global"
  13. )
  14. type testMQTTMessage struct {
  15. payload []byte
  16. }
  17. func (m testMQTTMessage) Duplicate() bool { return false }
  18. func (m testMQTTMessage) Qos() byte { return 1 }
  19. func (m testMQTTMessage) Retained() bool { return false }
  20. func (m testMQTTMessage) Topic() string { return "parking/test/gate/state" }
  21. func (m testMQTTMessage) MessageID() uint16 { return 1 }
  22. func (m testMQTTMessage) Payload() []byte { return m.payload }
  23. func (m testMQTTMessage) Ack() {}
  24. var _ mqtt.Message = testMQTTMessage{}
  25. func TestMQTTGateControllerClosedStateRemainsConnected(t *testing.T) {
  26. controller := NewMQTTGateController(nil)
  27. controller.onMessage(nil, testMQTTMessage{payload: []byte(`{"schema":"gate.state.v1","device_code":"gate-1","state":"closed"}`)})
  28. if !controller.IsGateConnected("gate-1") {
  29. t.Fatal("closed 状态应表示设备在线")
  30. }
  31. controller.mu.RLock()
  32. state := controller.gateState["gate-1"]
  33. controller.mu.RUnlock()
  34. if state != "closed" {
  35. t.Fatalf("闸杆状态未保存,得到 %q", state)
  36. }
  37. }
  38. func TestMQTTGateControllerRejectsAckFromOtherDevice(t *testing.T) {
  39. controller := NewMQTTGateController(nil)
  40. ackCh := make(chan gateAck, 1)
  41. controller.pending["cmd-1"] = pendingGateCommand{deviceCode: "gate-1", ackCh: ackCh}
  42. controller.onMessage(nil, testMQTTMessage{payload: []byte(`{"schema":"gate.ack.v1","cmd_id":"cmd-1","device_code":"gate-2","result":"success"}`)})
  43. select {
  44. case <-ackCh:
  45. t.Fatal("其他设备的 ACK 不应完成当前命令")
  46. default:
  47. }
  48. }
  49. func setupMQTTGateStatusTest(t *testing.T) {
  50. t.Helper()
  51. dsn := fmt.Sprintf("file:mqtt-gate-status-%d?mode=memory&cache=shared", time.Now().UnixNano())
  52. db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{})
  53. require.NoError(t, err)
  54. sqlDB, err := db.DB()
  55. require.NoError(t, err)
  56. t.Cleanup(func() {
  57. _ = sqlDB.Close()
  58. if global.GVA_DB == db {
  59. global.GVA_DB = nil
  60. }
  61. })
  62. global.GVA_DB = db
  63. global.GVA_LOG = zap.NewNop()
  64. t.Cleanup(func() { global.GVA_LOG = nil })
  65. require.NoError(t, db.AutoMigrate(&dao.UHFReader{}))
  66. }
  67. func TestMQTTGateControllerPersistsConfiguredDeviceStatus(t *testing.T) {
  68. setupMQTTGateStatusTest(t)
  69. require.NoError(t, global.GVA_DB.Create(&dao.UHFReader{
  70. DeviceCode: "GATE-01", DeviceName: "东门道闸", DeviceType: "gate",
  71. ConnectType: dao.ConnectTypeMQTT, ProvisionStatus: "provisioned",
  72. }).Error)
  73. controller := NewMQTTGateController(nil)
  74. controller.onMessage(nil, testMQTTMessage{payload: []byte(`{"schema":"gate.state.v1","device_code":"GATE-01","state":"open"}`)})
  75. var reader dao.UHFReader
  76. require.NoError(t, global.GVA_DB.Where("device_code = ?", "GATE-01").First(&reader).Error)
  77. require.Equal(t, "online", reader.Status)
  78. require.Equal(t, "online", reader.ProvisionStatus)
  79. require.NotNil(t, reader.LastOnlineTime)
  80. }
  81. func TestMQTTGateControllerMarksOnlyActiveMQTTDevicesOffline(t *testing.T) {
  82. setupMQTTGateStatusTest(t)
  83. readers := []dao.UHFReader{
  84. {DeviceCode: "MQTT-ONLINE", DeviceName: "MQTT 在线设备", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, IsActive: true, Status: "online"},
  85. {DeviceCode: "MQTT-DISABLED", DeviceName: "MQTT 停用设备", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, IsActive: false, Status: "online"},
  86. {DeviceCode: "TCP-ONLINE", DeviceName: "TCP 在线设备", DeviceType: "reader", ConnectType: dao.ConnectTypeTCP, IsActive: true, Status: "online"},
  87. }
  88. require.NoError(t, global.GVA_DB.Create(&readers).Error)
  89. require.NoError(t, global.GVA_DB.Model(&dao.UHFReader{}).
  90. Where("device_code = ?", "MQTT-DISABLED").
  91. Update("is_active", false).Error)
  92. controller := NewMQTTGateController(nil)
  93. require.NoError(t, controller.MarkConfiguredDevicesOffline())
  94. statuses := map[string]string{}
  95. var after []dao.UHFReader
  96. require.NoError(t, global.GVA_DB.Find(&after).Error)
  97. for _, reader := range after {
  98. statuses[reader.DeviceCode] = reader.Status
  99. }
  100. require.Equal(t, "offline", statuses["MQTT-ONLINE"])
  101. require.Equal(t, "online", statuses["MQTT-DISABLED"])
  102. require.Equal(t, "online", statuses["TCP-ONLINE"])
  103. }