# 设备接入层设计(摄像头 / 道闸 · MQTT + HTTP) > 状态:已落地(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 回调的设备,继电器/串口直连做挂在读写器上的道闸开闸。** --- ## 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-mqtt`(`internal/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) ```json { "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) ```json { "schema": "gate.state.v1", "device_code": "GATE-A01", "state": "open", "ts": 1730000000 } ``` `state` 取值:`open` / `closed` / `moving`。retained 消息保证系统重启后能立即恢复各道闸当前状态。 ### 3.3 开闸指令(gate/cmd) ```json { "schema": "gate.cmd.v1", "cmd_id": "3f2a1b9c-...", "device_code": "GATE-A01", "action": "open", "valid_time": 2 } ``` ### 3.4 指令回执(gate/cmd/ack) ```json { "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` 并存。 ```go 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` 订阅相机识别事件,统一桥接到现有进出场链路、通道事件、异常与前端推送。 ```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)。 --- ## 6. 启动装配(伪代码) ```go // 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` 新增) ```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.RecordIncident`(`device_offline`) | 道闸 LWT 离线埋点,复用既有去重 | | `incident.ResolveDeviceOffline` | 道闸恢复在线时关闭离线异常 | | `plugin/ws` | 识别事件经 ws 推给工作台 `entryExit.vue` | --- ## 9. HTTP 兜底通道(只支持 HTTP 回调的设备) 部分国产 LPR 相机只支持「HTTP POST 回调」,为其新增一个接收端点,内部转成同一条链路: ```go // 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.go`、`device_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 + ts` 或 `image_path` 关联。