低空智控平台 后端go
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

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