设备接入层设计.md 18 KB

设备接入层设计(摄像头 / 道闸 · MQTT + HTTP)

状态:已落地(MQTT 客户端 + 内嵌 broker 均已实现并编译通过) 落地源码:internal/service/mqttbroker/broker.gointernal/service/parking/mqtt_gate_controller.gointernal/service/parking/routing_gate_controller.gointernal/service/devicebus/device_bus.gointernal/initialize/mqtt.go 配套文档:doc/多岗亭改造方案.md 结论:MQTT 做消息总线(识别事件 / 状态 / 在线心跳 / 独立道闸指令),HTTP 传图片与兜底只支持 HTTP 回调的设备,继电器/串口直连做挂在读写器上的道闸开闸。


1. 总体架构

┌────────────┐  MQTT(event)  ┌────────────────────────────────────┐
│ LPR 相机    │──────────────▶│  DeviceBus(订阅识别事件)          │
│ (车牌识别)  │               │    │                               │
└────────────┘               │    ├─ 构造 PassageRequest           │
                             │    ├─ PassageService.HandlePassage  │──▶ 进场/出场 + 计费 + 数字票
┌────────────┐  HTTP(图片)   │    ├─ PushChannelEvent + 落库        │──▶ 前端工作台(ws)
│ 相机抓拍    │──────────────▶│    └─ 异常埋点(incident)            │
└────────────┘               └────────────────────────────────────┘

┌────────────┐  MQTT(state/LWT)  ┌──────────────────────────────┐
│ 道闸控制器  │──────────────────▶│  MQTTGateController           │
│ (独立网络)  │                   │    ├─ 订阅 state / ack / lwt   │
│            │◀──────────────────│    ├─ IsGateConnected()        │
│            │  MQTT(cmd, QoS1)  │    └─ 指令超时/失败 → incident  │
└────────────┘                   └──────────────────────────────┘

┌────────────┐  继电器/串口/TCP 直连(保持现状)
│ UHF 读写器  │──────────────────▶ 现有 Reader.CloseRelay1/2 + GateController
└────────────┘
  • MQTT Broker:已实现内嵌 mochi-mqttinternal/service/mqttbroker,随进程启动,无需额外进程);也可 embed-broker=false 连接外部 mosquitto
  • MQTT Client:Go 侧用 github.com/eclipse/paho.mqtt.golang(事实标准,需新增依赖)。

2. Topic 规划(按岗亭/通道分区,呼应多岗亭改造)

统一前缀 parking/,路由坐标 lot/{lot_id}/booth/{booth_id}/channel/{channel_id}

方向 Topic QoS Retained 用途
相机 → 系统 parking/lot/{lot}/booth/{booth}/channel/{ch}/camera/event 1 车牌识别事件
道闸 → 系统 parking/lot/{lot}/booth/{booth}/channel/{ch}/gate/state 1 道闸状态(开/关/在线)
系统 → 道闸 parking/lot/{lot}/booth/{booth}/channel/{ch}/gate/cmd 1 开/关闸指令
道闸 → 系统 parking/lot/{lot}/booth/{booth}/channel/{ch}/gate/cmd/ack 1 指令执行回执
道闸 → 系统 parking/lot/{lot}/booth/{booth}/channel/{ch}/gate/lwt 1 遗嘱消息(离线通知)

说明:GateController 接口只有 deviceCode,系统需在启动时从 DB 反查 UHFReader → Channel → Booth → ParkingLot 得到路由坐标(见 §4 LoadRoutes)。


3. 消息结构(JSON)

3.1 识别事件(camera/event)

{
  "schema": "camera.event.v1",
  "device_code": "CAM-A01",
  "plate_number": "粤A12345",
  "rfid_tag": "E28068940000500ABCDEF",
  "direction": "in",
  "ts": 1730000000,
  "image_path": "/capture/1730000000_A01.jpg",
  "confidence": 0.98
}

3.2 道闸状态(gate/state,retained)

{
  "schema": "gate.state.v1",
  "device_code": "GATE-A01",
  "state": "open",
  "ts": 1730000000
}

state 取值:open / closed / moving。retained 消息保证系统重启后能立即恢复各道闸当前状态。

3.3 开闸指令(gate/cmd)

{
  "schema": "gate.cmd.v1",
  "cmd_id": "3f2a1b9c-...",
  "device_code": "GATE-A01",
  "action": "open",
  "valid_time": 2
}

3.4 指令回执(gate/cmd/ack)

{
  "schema": "gate.ack.v1",
  "cmd_id": "3f2a1b9c-...",
  "device_code": "GATE-A01",
  "result": "success",
  "error": ""
}

4. 代码草案 A:internal/service/parking/mqtt_gate_controller.go

实现现有 GateController 接口,可经 parking.SetGateController(...) 注入,与当前 SimulatedGateController 并存。

package parking

import (
	"encoding/json"
	"errors"
	"fmt"
	"sync"
	"time"

	mqtt "github.com/eclipse/paho.mqtt.golang"
	"github.com/google/uuid"
	"wails-app/internal/dao"
	"wails-app/internal/global"
)

const (
	mqttQoS        = 1
	mqttCmdTimeout = 5 * time.Second
)

// gateTopicRoute 设备在 MQTT 主题树中的路由坐标。
type gateTopicRoute struct {
	DeviceCode string
	LotID      uint
	BoothID    uint
	ChannelID  uint
}

func (r gateTopicRoute) cmdTopic() string {
	return fmt.Sprintf("parking/lot/%d/booth/%d/channel/%d/gate/cmd", r.LotID, r.BoothID, r.ChannelID)
}

// MQTTGateController 通过 MQTT 下发道闸指令并等待设备 ack。
type MQTTGateController struct {
	client mqtt.Client

	mu        sync.RWMutex
	routes    map[string]gateTopicRoute // device_code -> route
	gateState map[string]bool           // device_code -> 最近上报 online/open
	pending   map[string]chan gateAck   // cmd_id -> ack channel
}

type gateCmd struct {
	Schema     string `json:"schema"`
	CmdID      string `json:"cmd_id"`
	DeviceCode string `json:"device_code"`
	Action     string `json:"action"`
	ValidTime  byte   `json:"valid_time"`
}

type gateAck struct {
	Schema     string `json:"schema"`
	CmdID      string `json:"cmd_id"`
	DeviceCode string `json:"device_code"`
	Result     string `json:"result"` // success / failed
	Error      string `json:"error"`
}

func NewMQTTGateController(client mqtt.Client) *MQTTGateController {
	return &MQTTGateController{
		client:    client,
		routes:    make(map[string]gateTopicRoute),
		gateState: make(map[string]bool),
		pending:   make(map[string]chan gateAck),
	}
}

// LoadRoutes 从 DB 反查 device_code -> lot/booth/channel,用于构造主题。
func (c *MQTTGateController) LoadRoutes() error {
	type row struct {
		DeviceCode   string
		ChannelID    uint
		BoothID      uint
		ParkingLotID uint
	}
	var rows []row
	err := global.GVA_DB.Model(&dao.UHFReader{}).
		Select("uhf_reader.device_code, uhf_reader.channel_id, channel.booth_id, channel.parking_lot_id").
		Joins("LEFT JOIN channel ON channel.id = uhf_reader.channel_id").
		Where("uhf_reader.channel_id > 0").
		Scan(&rows).Error
	if err != nil {
		return err
	}
	c.mu.Lock()
	defer c.mu.Unlock()
	for _, r := range rows {
		c.routes[r.DeviceCode] = gateTopicRoute{
			DeviceCode: r.DeviceCode, ChannelID: r.ChannelID,
			BoothID: r.BoothID, LotID: r.ParkingLotID,
		}
	}
	return nil
}

func (c *MQTTGateController) routeOf(deviceCode string) (gateTopicRoute, bool) {
	c.mu.RLock()
	defer c.mu.RUnlock()
	r, ok := c.routes[deviceCode]
	return r, ok
}

// OpenGate 实现 GateController 接口。
func (c *MQTTGateController) OpenGate(deviceCode string, validTime byte) error {
	return c.run(deviceCode, "open", validTime)
}

// CloseGate 实现 GateController 接口。
func (c *MQTTGateController) CloseGate(deviceCode string, validTime byte) error {
	return c.run(deviceCode, "close", validTime)
}

func (c *MQTTGateController) run(deviceCode, action string, validTime byte) error {
	route, ok := c.routeOf(deviceCode)
	if !ok {
		return fmt.Errorf("设备未配置 MQTT 路由: %s", deviceCode)
	}
	cmdID := uuid.NewString()
	ackCh := make(chan gateAck, 1)
	c.mu.Lock()
	c.pending[cmdID] = ackCh
	c.mu.Unlock()
	defer func() {
		c.mu.Lock()
		delete(c.pending, cmdID)
		c.mu.Unlock()
	}()

	payload, _ := json.Marshal(gateCmd{
		Schema: "gate.cmd.v1", CmdID: cmdID, DeviceCode: deviceCode,
		Action: action, ValidTime: validTime,
	})
	token := c.client.Publish(route.cmdTopic(), mqttQoS, false, payload)
	if token.Wait() && token.Error() != nil {
		return fmt.Errorf("MQTT 指令下发失败: %w", token.Error())
	}

	select {
	case ack := <-ackCh:
		if ack.Result != "success" {
			return fmt.Errorf("道闸执行失败: %s", ack.Error)
		}
		return nil
	case <-time.After(mqttCmdTimeout):
		return errors.New("道闸指令超时未收到 ack")
	}
}

// IsGateConnected 实现 GateController 接口(以最近 state/lwt 为准)。
func (c *MQTTGateController) IsGateConnected(deviceCode string) bool {
	c.mu.RLock()
	defer c.mu.RUnlock()
	return c.gateState[deviceCode]
}

// Subscribe 订阅 state / ack / lwt 主题。
func (c *MQTTGateController) Subscribe() error {
	for _, topic := range []string{
		"parking/lot/+/booth/+/channel/+/gate/state",
		"parking/lot/+/booth/+/channel/+/gate/cmd/ack",
		"parking/lot/+/booth/+/channel/+/gate/lwt",
	} {
		token := c.client.Subscribe(topic, mqttQoS, c.onMessage)
		if token.Wait() && token.Error() != nil {
			return token.Error()
		}
	}
	return nil
}

// onMessage 分派 state / ack / lwt。
func (c *MQTTGateController) onMessage(_ mqtt.Client, msg mqtt.Message) {
	payload := msg.Payload()
	// 先尝试按 schema 解析;ack 携带 cmd_id,state/lwt 携带 device_code + state。
	var probe struct {
		Schema     string `json:"schema"`
		CmdID      string `json:"cmd_id"`
		DeviceCode string `json:"device_code"`
		State      string `json:"state"`
	}
	_ = json.Unmarshal(payload, &probe)

	switch probe.Schema {
	case "gate.ack.v1":
		var ack gateAck
		if json.Unmarshal(payload, &ack) == nil {
			c.mu.RLock()
			ch := c.pending[ack.CmdID]
			c.mu.RUnlock()
			if ch != nil {
				ch <- ack
			}
		}
	case "gate.state.v1":
		c.mu.Lock()
		c.gateState[probe.DeviceCode] = probe.State == "open" || probe.State == "moving"
		c.mu.Unlock()
	case "gate.lwt.v1", "gate.lwt":
		// 遗嘱消息:设备离线。state 取 "offline"。
		c.mu.Lock()
		c.gateState[probe.DeviceCode] = false
		c.mu.Unlock()
	}
}

5. 代码草案 B:internal/service/parking/device_bus.go

订阅相机识别事件,统一桥接到现有进出场链路、通道事件、异常与前端推送。

package parking

import (
	"encoding/json"
	"fmt"

	mqtt "github.com/eclipse/paho.mqtt.golang"
	"wails-app/internal/global"
	common "wails-app/internal/model/common"
	incidentService "wails-app/internal/modules/incident/service"
	"wails-app/internal/service/uhf"
)

// DeviceBus 订阅相机/道闸的 MQTT 消息,桥接到业务层。
type DeviceBus struct {
	passage    *PassageService
	incident   *incidentService.IncidentService
}

func NewDeviceBus() *DeviceBus {
	return &DeviceBus{
		passage:  &PassageService{},
		incident: incidentService.NewIncidentService(),
	}
}

type cameraEvent struct {
	Schema      string  `json:"schema"`
	DeviceCode  string  `json:"device_code"`
	PlateNumber string  `json:"plate_number"`
	RFIDTag     string  `json:"rfid_tag"`
	Direction   string  `json:"direction"`
	TS          int64   `json:"ts"`
	ImagePath   string  `json:"image_path"`
	Confidence  float64 `json:"confidence"`
}

// Subscribe 订阅识别事件与道闸 LWT。
func (b *DeviceBus) Subscribe(client mqtt.Client) error {
	for _, topic := range []string{
		"parking/lot/+/booth/+/channel/+/camera/event",
		"parking/lot/+/booth/+/channel/+/gate/lwt",
	} {
		token := client.Subscribe(topic, 1, b.onMessage)
		if token.Wait() && token.Error() != nil {
			return token.Error()
		}
	}
	return nil
}

func (b *DeviceBus) onMessage(_ mqtt.Client, msg mqtt.Message) {
	payload := msg.Payload()
	var probe struct {
		Schema string `json:"schema"`
	}
	_ = json.Unmarshal(payload, &probe)

	switch probe.Schema {
	case "camera.event.v1":
		b.onCameraEvent(payload)
	case "gate.lwt.v1", "gate.lwt":
		b.onGateOffline(payload)
	}
}

func (b *DeviceBus) onCameraEvent(payload []byte) {
	var ev cameraEvent
	if err := json.Unmarshal(payload, &ev); err != nil {
		return
	}

	req := common.PassageRequest{
		DeviceCode:    ev.DeviceCode,
		PlateNumber:   ev.PlateNumber,
		RFIDTag:       ev.RFIDTag,
		Direction:     ev.Direction,
		Image:         ev.ImagePath,
		TriggerSource: "camera",
	}

	// 复用统一进出场入口:黑名单 / 临时车规则 / 满位 / 计费 全在这里。
	result, err := b.passage.HandlePassage(req)

	status := "completed"
	message := "识别通过"
	fee := 0.0
	stayTime := int64(0)
	if err != nil {
		status = "failed"
		message = err.Error()
	} else if result != nil {
		fee = result.Fee
		stayTime = result.StayTime
		if result.Message != "" {
			message = result.Message
		}
	}

	// 1) 推送到现有通道事件队列(工作台实时刷新),并可按岗亭/区域过滤。
	//    多岗亭改造后建议在此同时落库 passage_event。
	uhf.PushChannelEvent(uhf.ChannelEvent{
		DeviceCode:  ev.DeviceCode,
		RFIDTag:     ev.RFIDTag,
		PlateNumber: ev.PlateNumber,
		Direction:   ev.Direction,
		Timestamp:   ev.TS,
		Status:      status,
		Message:     message,
		Fee:         fee,
		StayTime:    stayTime,
	})
}

// onGateOffline 道闸遗嘱消息 → 设备离线异常(对接 incident 既有去重)。
func (b *DeviceBus) onGateOffline(payload []byte) {
	var probe struct {
		DeviceCode string `json:"device_code"`
	}
	if json.Unmarshal(payload, &probe) != nil || probe.DeviceCode == "" {
		return
	}
	b.incident.RecordIncident(incidentService.RecordIncidentRequest{
		Category:    incidentService.CategoryDeviceOffline,
		Source:      incidentService.SourceDevice,
		DeviceCode:  probe.DeviceCode,
		Description: "道闸 MQTT 遗嘱触发,设备离线",
	})
}

补充:设备恢复在线时,由 gate/state(retained)上报触发,调用 incident.ResolveDeviceOffline(deviceCode) 关闭未处理的离线异常;MQTTGateController.onMessagegate.state.v1 分支可顺带调用一次(best-effort)。


6. 启动装配(伪代码)

// internal/initialize/mqtt.go(新增)
func InitDeviceBus() {
	opts := mqtt.NewClientOptions().
		AddBroker(global.GVA_CONFIG.Mqtt.Broker).
		SetClientID("smart-parking").
		SetWill("parking/device/lwt", `{"device_code":"smart-parking","state":"offline"}`, 1, true).
		SetAutoReconnect(true)

	client := mqtt.NewClient(opts)
	if token := client.Connect(); token.Wait() && token.Error() != nil {
		global.GVA_LOG.Error("MQTT 连接失败", zap.Error(token.Error()))
		return
	}

	gate := parking.NewMQTTGateController(client)
	_ = gate.LoadRoutes()
	_ = gate.Subscribe()
	parking.SetGateController(gate) // 与现有 SimulatedGateController 互斥注入

	bus := parking.NewDeviceBus()
	_ = bus.Subscribe(client)
}

7. 配置项(config.yaml 新增)

mqtt:
  enabled: false                    # 总开关
  embed-broker: true                # 是否启动内嵌 mochi-mqtt broker
  broker-listen: ":1883"            # 内嵌 broker 监听地址
  broker: "tcp://127.0.0.1:1883"    # 客户端连接地址(内嵌/外部 broker 均连此地址)
  client-id: "smart-parking"
  username: ""
  password: ""

对应 internal/config/ 新增 mqtt.go 结构体,并挂到 config.yaml 加载。


8. 与现有代码的对接点

现有代码 对接方式
GateController 接口(service/parking/gate.go MQTTGateController 直接实现,SetGateController 注入
GetGateRuntimeStatus / isSimulator() MQTT 控制器不实现 isSimulator()Simulated=false
PassageService.HandlePassage 相机事件统一走此入口,黑名单/临时车/满位/计费复用
uhf.PushChannelEvent / GetChannelEvents 识别结果推送;多岗亭改造后补落库 + 按岗亭过滤
incident.RecordIncidentdevice_offline 道闸 LWT 离线埋点,复用既有去重
incident.ResolveDeviceOffline 道闸恢复在线时关闭离线异常
plugin/ws 识别事件经 ws 推给工作台 entryExit.vue

9. HTTP 兜底通道(只支持 HTTP 回调的设备)

部分国产 LPR 相机只支持「HTTP POST 回调」,为其新增一个接收端点,内部转成同一条链路:

// internal/api/v1/parking/camera_http.go(新增)
func CameraEventCallback(c *gin.Context) {
	var ev cameraEvent // 复用 §5 结构
	if err := c.ShouldBindJSON(&ev); err != nil {
		response.FailWithMessage(err.Error(), c)
		return
	}
	// 转成统一 PassageRequest 走 DeviceBus 同一条链路
	// (图片本体走 multipart 上传或 image_url 由系统 HTTP 拉取)
	response.OkWithMessage("ok", c)
}

图片本体一律走 HTTP(POST multipart 上传或按 image_path/URL 拉取),不塞进 MQTT,避免撑爆 broker。


10. 依赖与实施步骤

  1. 依赖已引入:github.com/eclipse/paho.mqtt.golang(客户端)+ github.com/mochi-mqtt/server/v2(内嵌 broker)。
  2. 新增 internal/config/mqtt.go + config.yaml 配置。
  3. 新增 internal/service/parking/mqtt_gate_controller.godevice_bus.go
  4. 新增 internal/initialize/mqtt.go,在 InitBackend/seed 流程调用 InitDeviceBus()
  5. 相机/道闸厂商按 §2、§3 的 topic 与 schema 接入。
  6. 联调顺序:模拟道闸 → 相机事件 → LWT 离线 → 独立网络道闸指令回执。

11. 落地边界(避免过度设计)

  • 单进程、局域网、设备数量有限:不引入 Kafka/时序库,MQTT broker 已内嵌 mochi-mqtt,随应用启动。
  • 开闸硬实时优先:挂在读写器继电器上的道闸继续走串口/TCP 直连,MQTT 只做「独立网络道闸」和「状态/审计上报」,不强制改造现有硬件链路。
  • 图片不进 MQTT:元数据走 MQTT、图片走 HTTP,两者用 device_code + tsimage_path 关联。