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) }