You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
408 lines
13 KiB
408 lines
13 KiB
package service
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"strconv"
|
|
"strings"
|
|
|
|
paho "github.com/eclipse/paho.mqtt.golang"
|
|
|
|
"laic-backend/cache"
|
|
"laic-backend/common"
|
|
"laic-backend/logger"
|
|
"laic-backend/mqtt"
|
|
"laic-backend/websocket"
|
|
)
|
|
|
|
// 上行消息业务载荷结构(协议文档 §7,均位于通用外层 payload 内)
|
|
|
|
// statusOnline §7.1 工控机在线状态
|
|
type statusOnline struct {
|
|
Status string `json:"status"` // online / offline / degraded / updating
|
|
BootID string `json:"bootId"`
|
|
DockIDSource string `json:"dockIdSource"`
|
|
SoftwareVersion string `json:"softwareVersion"`
|
|
ProtocolVersion string `json:"protocolVersion"`
|
|
UptimeSec int64 `json:"uptimeSec"`
|
|
TimeSynced bool `json:"timeSynced"`
|
|
MqttConnected bool `json:"mqttConnected"`
|
|
ModbusConnected bool `json:"modbusConnected"`
|
|
MavlinkConnected bool `json:"mavlinkConnected"`
|
|
Updating bool `json:"updating"`
|
|
Name string `json:"name"`
|
|
Location string `json:"location"`
|
|
Latitude float64 `json:"latitude"`
|
|
Longitude float64 `json:"longitude"`
|
|
}
|
|
|
|
// dockState §7.2 机巢状态(仅解析业务关心字段,其余经 flatten 存 Redis)
|
|
type dockState struct {
|
|
PlcConnected bool `json:"plcConnected"`
|
|
ControlMode string `json:"controlMode"`
|
|
ChargingState string `json:"chargingState"`
|
|
EmergencyStop bool `json:"emergencyStop"`
|
|
AlarmCodes []string `json:"alarmCodes"`
|
|
DronePresent *bool `json:"dronePresent"`
|
|
}
|
|
|
|
// droneState §7.3 无人机状态
|
|
type droneState struct {
|
|
DroneSN string `json:"droneSn"`
|
|
Name string `json:"name"`
|
|
CurrentSysID int `json:"currentSysId"`
|
|
Online bool `json:"online"`
|
|
Armed bool `json:"armed"`
|
|
FlightMode string `json:"flightMode"`
|
|
FlightModeCode int `json:"flightModeCode"`
|
|
Latitude float64 `json:"latitude"`
|
|
Longitude float64 `json:"longitude"`
|
|
Altitude float64 `json:"altitude"`
|
|
GroundSpeed float64 `json:"groundSpeed"`
|
|
Roll float64 `json:"roll"`
|
|
Pitch float64 `json:"pitch"`
|
|
Yaw float64 `json:"yaw"`
|
|
BatteryPercent float64 `json:"batteryPercent"`
|
|
BatteryVoltage float64 `json:"batteryVoltage"`
|
|
BatteryCurrent float64 `json:"batteryCurrent"`
|
|
Satellites int `json:"satellites"`
|
|
GpsQuality string `json:"gpsQuality"`
|
|
LinkQuality float64 `json:"linkQuality"`
|
|
HomeSet bool `json:"homeSet"`
|
|
AlarmCodes []string `json:"alarmCodes"`
|
|
}
|
|
|
|
// commandAck §5.3 指令应答
|
|
type commandAck struct {
|
|
CommandID string `json:"commandId"`
|
|
Accepted bool `json:"accepted"`
|
|
ResultCode string `json:"resultCode"`
|
|
}
|
|
|
|
// otaReported §8 OTA 升级进度上报
|
|
type otaReported struct {
|
|
UpdateID string `json:"updateId"`
|
|
Component string `json:"component"`
|
|
TargetVersion string `json:"targetVersion"`
|
|
CurrentVersion string `json:"currentVersion"`
|
|
Status string `json:"status"`
|
|
Progress int `json:"progress"`
|
|
Message string `json:"message"`
|
|
ErrorCode *string `json:"errorCode"`
|
|
ReportedAt int64 `json:"reportedAt"`
|
|
}
|
|
|
|
// videoState §7.8 视频推流状态
|
|
type videoState struct {
|
|
Version int64 `json:"version"`
|
|
EventID string `json:"eventId"`
|
|
Provider string `json:"provider"`
|
|
Phase string `json:"phase"`
|
|
InputOnline bool `json:"inputOnline"`
|
|
Streaming bool `json:"streaming"`
|
|
InputCodec string `json:"inputCodec"`
|
|
UplinkProtocol string `json:"uplinkProtocol"`
|
|
StreamSessionID string `json:"streamSessionId"`
|
|
Width *int `json:"width"`
|
|
Height *int `json:"height"`
|
|
FrameRate *float64 `json:"frameRate"`
|
|
BitrateBps uint64 `json:"bitrateBps"`
|
|
RetryCount int `json:"retryCount"`
|
|
StopReason *string `json:"stopReason"`
|
|
ErrorCode *string `json:"errorCode"`
|
|
UpdatedAt int64 `json:"updatedAt"`
|
|
}
|
|
|
|
// InitMQTTSubscriber 注册 MQTT 订阅处理器
|
|
func InitMQTTSubscriber() {
|
|
mqtt.Subscribe(map[string]byte{
|
|
"dock-edge/v1/dock/+/status/online": 1,
|
|
"dock-edge/v1/dock/+/state/dock": 1,
|
|
"dock-edge/v1/dock/+/state/drone": 1,
|
|
"dock-edge/v1/dock/+/state/workflow": 1,
|
|
"dock-edge/v1/dock/+/state/video": 1,
|
|
"dock-edge/v1/dock/+/internal/original-video": 1,
|
|
"dock-edge/v1/dock/+/telemetry": 0,
|
|
"dock-edge/v1/dock/+/command/ack": 1,
|
|
"dock-edge/v1/dock/+/ota/reported": 1,
|
|
}, routeMessage)
|
|
}
|
|
|
|
// routeMessage 按 topic 分发到对应处理器
|
|
// topic 结构:dock-edge/v1/dock/{dockId}/{category}/{sub}
|
|
// 消息体为通用外层 Envelope,业务字段在 payload 内。
|
|
func routeMessage(_ paho.Client, msg paho.Message) {
|
|
topic := msg.Topic()
|
|
parts := strings.Split(topic, "/")
|
|
if len(parts) < 5 {
|
|
logger.WARN("无法解析 MQTT topic:", topic)
|
|
return
|
|
}
|
|
topicDockID := parts[3]
|
|
dockID := topicDockID
|
|
category := parts[4]
|
|
sub := ""
|
|
if len(parts) > 5 {
|
|
sub = parts[5]
|
|
}
|
|
|
|
env := mqtt.ParseEnvelope(msg.Payload())
|
|
if env.DockID != "" && env.DockID != topicDockID {
|
|
logger.WARN("MQTT envelope dockId 与 topic 不一致", topicDockID, env.DockID)
|
|
return
|
|
}
|
|
if env.DockID != "" {
|
|
dockID = env.DockID
|
|
}
|
|
|
|
switch category {
|
|
case "status":
|
|
if sub == "online" {
|
|
handleStatus(dockID, env)
|
|
}
|
|
case "state":
|
|
switch sub {
|
|
case "dock":
|
|
handleStateDock(dockID, env)
|
|
case "drone":
|
|
handleStateDrone(dockID, env)
|
|
case "workflow":
|
|
handleWorkflow(dockID, env)
|
|
case "video":
|
|
handleVideo(dockID, env)
|
|
}
|
|
case "telemetry":
|
|
handleTelemetry(dockID, env)
|
|
case "command":
|
|
if sub == "ack" {
|
|
handleCommandAck(dockID, env)
|
|
}
|
|
case "ota":
|
|
if sub == "reported" {
|
|
handleOtaReported(dockID, env)
|
|
}
|
|
case "internal":
|
|
if sub == "original-video" {
|
|
handleOriginalVideo(dockID, env)
|
|
}
|
|
default:
|
|
logger.DEBUG("未知 topic:", topic)
|
|
}
|
|
}
|
|
|
|
// handleStatus 在线/离线 → 设备自发现 + 在线集合 + 软件版本
|
|
func handleStatus(dockID string, env *mqtt.Envelope) {
|
|
var st statusOnline
|
|
if err := json.Unmarshal(env.Payload, &st); err != nil {
|
|
logger.ERROR("解析 status/online 失败", err)
|
|
return
|
|
}
|
|
DefaultDockService.EnsureFromOnline(dockID, st)
|
|
if st.Status == "online" {
|
|
refreshDeviceHeartbeat(cache.DockHeartbeatKeyOf(dockID))
|
|
}
|
|
}
|
|
|
|
// handleStateDock 机巢状态 → Redis Hash + 告警 diff
|
|
func handleStateDock(dockID string, env *mqtt.Envelope) {
|
|
var raw map[string]any
|
|
_ = json.Unmarshal(env.Payload, &raw)
|
|
_ = common.HashSetValues(cache.DockStatusKeyOf(dockID), flatten(raw))
|
|
refreshDeviceHeartbeat(cache.DockHeartbeatKeyOf(dockID))
|
|
|
|
var st dockState
|
|
if err := json.Unmarshal(env.Payload, &st); err != nil {
|
|
logger.ERROR("解析 state/dock 失败", err)
|
|
return
|
|
}
|
|
DefaultAlarmService.SyncByState(dockID, "dock", st.AlarmCodes)
|
|
}
|
|
|
|
// handleStateDrone 无人机状态 → 自动关联 + Redis 遥测 + 告警 diff
|
|
func handleStateDrone(dockID string, env *mqtt.Envelope) {
|
|
var raw map[string]any
|
|
_ = json.Unmarshal(env.Payload, &raw)
|
|
_ = common.HashSetValues(cache.DroneTelemetryKeyOf(dockID), flatten(raw))
|
|
|
|
var st droneState
|
|
if err := json.Unmarshal(env.Payload, &st); err != nil {
|
|
logger.ERROR("解析 state/drone 失败", err)
|
|
return
|
|
}
|
|
droneSN := st.DroneSN
|
|
if droneSN == "" && env.DroneSN != nil {
|
|
droneSN = *env.DroneSN
|
|
}
|
|
battery := int(st.BatteryPercent)
|
|
DefaultDroneService.EnsureFromState(dockID, droneSN, st.Name, st.Online, battery, "")
|
|
if droneSN != "" && st.Online {
|
|
refreshDeviceHeartbeat(cache.DroneHeartbeatKeyOf(droneSN))
|
|
}
|
|
DefaultAlarmService.SyncByState(dockID, "drone", st.AlarmCodes)
|
|
DefaultTelemetryStore.SetDroneSN(dockID, droneSN)
|
|
}
|
|
|
|
// handleTelemetry 高频遥测 → Redis 最新值 + TDengine 批量写入 + 轨迹采集
|
|
func handleTelemetry(dockID string, env *mqtt.Envelope) {
|
|
var raw map[string]any
|
|
if err := json.Unmarshal(env.Payload, &raw); err != nil {
|
|
logger.ERROR("解析 telemetry 失败", err)
|
|
return
|
|
}
|
|
_ = common.HashSetValues(cache.DroneTelemetryKeyOf(dockID), flatten(raw))
|
|
|
|
var p TelemetryPoint
|
|
if err := json.Unmarshal(env.Payload, &p); err == nil {
|
|
DefaultTelemetryStore.Append(dockID, p)
|
|
DefaultTrajectoryStore.Append(dockID, p)
|
|
}
|
|
if env.DroneSN != nil && *env.DroneSN != "" {
|
|
DefaultTelemetryStore.SetDroneSN(dockID, *env.DroneSN)
|
|
refreshDeviceHeartbeat(cache.DroneHeartbeatKeyOf(*env.DroneSN))
|
|
}
|
|
}
|
|
|
|
// handleWorkflow 工作流状态 → Redis Hash + WS 推送 + 持久化 + 终态回写任务
|
|
func handleWorkflow(dockID string, env *mqtt.Envelope) {
|
|
var raw map[string]any
|
|
_ = json.Unmarshal(env.Payload, &raw)
|
|
_ = common.HashSetValues(cache.WorkflowKeyOf(dockID), flatten(raw))
|
|
|
|
var st WorkflowStateIn
|
|
if err := json.Unmarshal(env.Payload, &st); err != nil {
|
|
logger.ERROR("解析 state/workflow 失败", err)
|
|
return
|
|
}
|
|
DefaultWorkflowService.Upsert(dockID, env.RequestID, &st)
|
|
broadcast("workflow.state", map[string]any{
|
|
"dockId": dockID,
|
|
"commandId": st.CommandID,
|
|
"taskId": st.TaskID,
|
|
"missionId": st.MissionID,
|
|
"state": st.State,
|
|
"step": st.Step,
|
|
"resultCode": st.ResultCode,
|
|
})
|
|
}
|
|
|
|
// handleVideo 视频推流状态 → Redis Hash + WS 推送 + 直播会话联动
|
|
func handleVideo(dockID string, env *mqtt.Envelope) {
|
|
var raw map[string]any
|
|
_ = json.Unmarshal(env.Payload, &raw)
|
|
var st videoState
|
|
if err := json.Unmarshal(env.Payload, &st); err != nil {
|
|
logger.ERROR("解析 state/video 失败", err)
|
|
return
|
|
}
|
|
if st.EventID == "" {
|
|
st.EventID = env.EventID
|
|
}
|
|
if st.Version == 0 {
|
|
st.Version = env.Version
|
|
}
|
|
if st.StreamSessionID == "" || st.EventID == "" || st.Version <= 0 || st.UpdatedAt <= 0 {
|
|
logger.WARN("忽略缺少顺序元数据的视频状态", dockID, st.StreamSessionID)
|
|
return
|
|
}
|
|
if !DefaultLiveService.OnVideoState(dockID, st.Streaming, st.StreamSessionID, st.EventID, st.Version, st.UpdatedAt) {
|
|
return
|
|
}
|
|
_ = common.HashSetValues(cache.VideoKeyOf(dockID), flatten(raw))
|
|
broadcast("video.state", map[string]any{
|
|
"dockId": dockID,
|
|
"streaming": st.Streaming,
|
|
"streamSessionId": st.StreamSessionID,
|
|
"bitrate": st.BitrateBps,
|
|
"phase": st.Phase,
|
|
"errorCode": st.ErrorCode,
|
|
})
|
|
}
|
|
|
|
// handleOriginalVideo 处理后台与 Mock 的内部原始视频上传事件。
|
|
func handleOriginalVideo(dockID string, env *mqtt.Envelope) {
|
|
if env.DockID != "" && env.DockID != dockID {
|
|
logger.WARN("原始视频事件机巢不匹配", dockID, env.DockID)
|
|
return
|
|
}
|
|
var event OriginalVideoEvent
|
|
if err := json.Unmarshal(env.Payload, &event); err != nil {
|
|
logger.ERROR("解析原始视频事件失败", err)
|
|
return
|
|
}
|
|
if err := DefaultVideoService.CompleteOriginalVideo(dockID, &event); err != nil {
|
|
logger.ERROR("处理原始视频事件失败", err)
|
|
return
|
|
}
|
|
broadcast("video.original", map[string]any{
|
|
"dockId": dockID, "videoId": event.VideoID, "executionId": event.ExecutionID, "eventType": event.EventType,
|
|
})
|
|
}
|
|
|
|
// handleCommandAck 指令应答 → 更新 device_command_log
|
|
func handleCommandAck(dockID string, env *mqtt.Envelope) {
|
|
var ack commandAck
|
|
if err := json.Unmarshal(env.Payload, &ack); err != nil {
|
|
logger.ERROR("解析 command/ack 失败", err)
|
|
return
|
|
}
|
|
DefaultCommandService.HandleAck(dockID, env.RequestID, ack.CommandID, ack.Accepted, ack.ResultCode)
|
|
}
|
|
|
|
// handleOtaReported OTA 升级进度 → Redis 最新值 + WS 推送
|
|
func handleOtaReported(dockID string, env *mqtt.Envelope) {
|
|
var r otaReported
|
|
if err := json.Unmarshal(env.Payload, &r); err != nil {
|
|
logger.ERROR("解析 ota/reported 失败", err)
|
|
return
|
|
}
|
|
_ = common.HashSetValues(cache.OtaStatusKeyOf(dockID), map[string]any{
|
|
"updateId": r.UpdateID,
|
|
"component": r.Component,
|
|
"targetVersion": r.TargetVersion,
|
|
"currentVersion": r.CurrentVersion,
|
|
"status": r.Status,
|
|
"progress": strconv.Itoa(r.Progress),
|
|
"message": r.Message,
|
|
})
|
|
broadcast("ota.progress", map[string]any{
|
|
"dockId": dockID,
|
|
"updateId": r.UpdateID,
|
|
"targetVersion": r.TargetVersion,
|
|
"status": r.Status,
|
|
"progress": r.Progress,
|
|
"message": r.Message,
|
|
})
|
|
}
|
|
|
|
// flatten 将 JSON map 展平为 Redis Hash 可存储的 string 值
|
|
func flatten(data map[string]any) map[string]any {
|
|
out := make(map[string]any, len(data))
|
|
for k, v := range data {
|
|
switch val := v.(type) {
|
|
case string:
|
|
out[k] = val
|
|
case float64:
|
|
out[k] = strconv.FormatFloat(val, 'f', -1, 64)
|
|
case bool:
|
|
out[k] = strconv.FormatBool(val)
|
|
case nil:
|
|
out[k] = ""
|
|
default:
|
|
if b, err := json.Marshal(val); err == nil {
|
|
out[k] = string(b)
|
|
} else {
|
|
out[k] = fmt.Sprint(val)
|
|
}
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// broadcast 向有权访问该机巢的 WS 客户端推送脱敏事件。
|
|
func broadcast(eventType string, data map[string]any) {
|
|
dockID, _ := data["dockId"].(string)
|
|
if dockID == "" {
|
|
return
|
|
}
|
|
websocket.DefaultHub.BroadcastToDockJSON(dockID, eventType, data)
|
|
}
|
|
|