状态:已落地(MQTT 客户端 + 内嵌 broker 均已实现并编译通过) 落地源码:
internal/service/mqttbroker/broker.go、internal/service/parking/mqtt_gate_controller.go、internal/service/parking/routing_gate_controller.go、internal/service/devicebus/device_bus.go、internal/initialize/mqtt.go配套文档:doc/多岗亭改造方案.md结论:MQTT 做消息总线(识别事件 / 状态 / 在线心跳 / 独立道闸指令),HTTP 传图片与兜底只支持 HTTP 回调的设备,继电器/串口直连做挂在读写器上的道闸开闸。
┌────────────┐ 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
└────────────┘
mochi-mqtt(internal/service/mqttbroker,随进程启动,无需额外进程);也可 embed-broker=false 连接外部 mosquitto。github.com/eclipse/paho.mqtt.golang(事实标准,需新增依赖)。统一前缀 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得到路由坐标(见 §4LoadRoutes)。
{
"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
}
{
"schema": "gate.state.v1",
"device_code": "GATE-A01",
"state": "open",
"ts": 1730000000
}
state 取值:open / closed / moving。retained 消息保证系统重启后能立即恢复各道闸当前状态。
{
"schema": "gate.cmd.v1",
"cmd_id": "3f2a1b9c-...",
"device_code": "GATE-A01",
"action": "open",
"valid_time": 2
}
{
"schema": "gate.ack.v1",
"cmd_id": "3f2a1b9c-...",
"device_code": "GATE-A01",
"result": "success",
"error": ""
}
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()
}
}
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.onMessage的gate.state.v1分支可顺带调用一次(best-effort)。
// 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)
}
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 加载。
| 现有代码 | 对接方式 |
|---|---|
GateController 接口(service/parking/gate.go) |
MQTTGateController 直接实现,SetGateController 注入 |
GetGateRuntimeStatus / isSimulator() |
MQTT 控制器不实现 isSimulator(),Simulated=false |
PassageService.HandlePassage |
相机事件统一走此入口,黑名单/临时车/满位/计费复用 |
uhf.PushChannelEvent / GetChannelEvents |
识别结果推送;多岗亭改造后补落库 + 按岗亭过滤 |
incident.RecordIncident(device_offline) |
道闸 LWT 离线埋点,复用既有去重 |
incident.ResolveDeviceOffline |
道闸恢复在线时关闭离线异常 |
plugin/ws |
识别事件经 ws 推给工作台 entryExit.vue |
部分国产 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。
github.com/eclipse/paho.mqtt.golang(客户端)+ github.com/mochi-mqtt/server/v2(内嵌 broker)。internal/config/mqtt.go + config.yaml 配置。internal/service/parking/mqtt_gate_controller.go、device_bus.go。internal/initialize/mqtt.go,在 InitBackend/seed 流程调用 InitDeviceBus()。device_code + ts 或 image_path 关联。