package parking import ( "fmt" "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" ) type testMQTTMessage struct { payload []byte } func (m testMQTTMessage) Duplicate() bool { return false } func (m testMQTTMessage) Qos() byte { return 1 } func (m testMQTTMessage) Retained() bool { return false } func (m testMQTTMessage) Topic() string { return "parking/test/gate/state" } func (m testMQTTMessage) MessageID() uint16 { return 1 } func (m testMQTTMessage) Payload() []byte { return m.payload } func (m testMQTTMessage) Ack() {} var _ mqtt.Message = testMQTTMessage{} func TestMQTTGateControllerClosedStateRemainsConnected(t *testing.T) { controller := NewMQTTGateController(nil) controller.onMessage(nil, testMQTTMessage{payload: []byte(`{"schema":"gate.state.v1","device_code":"gate-1","state":"closed"}`)}) if !controller.IsGateConnected("gate-1") { t.Fatal("closed 状态应表示设备在线") } controller.mu.RLock() state := controller.gateState["gate-1"] controller.mu.RUnlock() if state != "closed" { t.Fatalf("闸杆状态未保存,得到 %q", state) } } func TestMQTTGateControllerRejectsAckFromOtherDevice(t *testing.T) { controller := NewMQTTGateController(nil) ackCh := make(chan gateAck, 1) controller.pending["cmd-1"] = pendingGateCommand{deviceCode: "gate-1", ackCh: ackCh} controller.onMessage(nil, testMQTTMessage{payload: []byte(`{"schema":"gate.ack.v1","cmd_id":"cmd-1","device_code":"gate-2","result":"success"}`)}) select { case <-ackCh: t.Fatal("其他设备的 ACK 不应完成当前命令") default: } } func setupMQTTGateStatusTest(t *testing.T) { t.Helper() dsn := fmt.Sprintf("file:mqtt-gate-status-%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) t.Cleanup(func() { _ = sqlDB.Close() if global.GVA_DB == db { global.GVA_DB = nil } }) global.GVA_DB = db global.GVA_LOG = zap.NewNop() t.Cleanup(func() { global.GVA_LOG = nil }) require.NoError(t, db.AutoMigrate(&dao.UHFReader{})) } func TestMQTTGateControllerPersistsConfiguredDeviceStatus(t *testing.T) { setupMQTTGateStatusTest(t) require.NoError(t, global.GVA_DB.Create(&dao.UHFReader{ DeviceCode: "GATE-01", DeviceName: "东门道闸", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, ProvisionStatus: "provisioned", }).Error) controller := NewMQTTGateController(nil) controller.onMessage(nil, testMQTTMessage{payload: []byte(`{"schema":"gate.state.v1","device_code":"GATE-01","state":"open"}`)}) var reader dao.UHFReader require.NoError(t, global.GVA_DB.Where("device_code = ?", "GATE-01").First(&reader).Error) require.Equal(t, "online", reader.Status) require.Equal(t, "online", reader.ProvisionStatus) require.NotNil(t, reader.LastOnlineTime) } func TestMQTTGateControllerMarksOnlyActiveMQTTDevicesOffline(t *testing.T) { setupMQTTGateStatusTest(t) readers := []dao.UHFReader{ {DeviceCode: "MQTT-ONLINE", DeviceName: "MQTT 在线设备", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, IsActive: true, Status: "online"}, {DeviceCode: "MQTT-DISABLED", DeviceName: "MQTT 停用设备", DeviceType: "gate", ConnectType: dao.ConnectTypeMQTT, IsActive: false, Status: "online"}, {DeviceCode: "TCP-ONLINE", DeviceName: "TCP 在线设备", DeviceType: "reader", ConnectType: dao.ConnectTypeTCP, IsActive: true, Status: "online"}, } require.NoError(t, global.GVA_DB.Create(&readers).Error) require.NoError(t, global.GVA_DB.Model(&dao.UHFReader{}). Where("device_code = ?", "MQTT-DISABLED"). Update("is_active", false).Error) controller := NewMQTTGateController(nil) require.NoError(t, controller.MarkConfiguredDevicesOffline()) statuses := map[string]string{} var after []dao.UHFReader require.NoError(t, global.GVA_DB.Find(&after).Error) for _, reader := range after { statuses[reader.DeviceCode] = reader.Status } require.Equal(t, "offline", statuses["MQTT-ONLINE"]) require.Equal(t, "online", statuses["MQTT-DISABLED"]) require.Equal(t, "online", statuses["TCP-ONLINE"]) }