Commit b1d18fde by 真的三个金的鑫

支持 Pilot2 上云:云控授权、DRC 设备 Broker 与 RC 飞控分支。

对 Pilot RC 走 cloud_control_auth;DRC 区分 Web WSS 与设备 MQTT;拒发无应答的 flighttask;增强枚举/OSD 容错。

Co-authored-by: Cursor <cursoragent@cursor.com>
parent 24356261
package com.dji.sdk.cloudapi.control;
import com.dji.sdk.common.BaseModel;
import javax.validation.constraints.NotNull;
import javax.validation.constraints.Size;
import java.util.List;
/**
* Pilot 云端申请飞行控制权:cloud_control_auth_request
* @see <a href="https://developer.dji.com/doc/cloud-api-tutorial/cn/api-reference/pilot-to-cloud/mqtt/rc-pro/drc.html">Pilot DRC</a>
*/
public class CloudControlAuthRequest extends BaseModel {
@NotNull
private String userId;
@NotNull
private String userCallsign;
@NotNull
@Size(min = 1)
private List<String> controlKeys;
public CloudControlAuthRequest() {
}
public String getUserId() {
return userId;
}
public CloudControlAuthRequest setUserId(String userId) {
this.userId = userId;
return this;
}
public String getUserCallsign() {
return userCallsign;
}
public CloudControlAuthRequest setUserCallsign(String userCallsign) {
this.userCallsign = userCallsign;
return this;
}
public List<String> getControlKeys() {
return controlKeys;
}
public CloudControlAuthRequest setControlKeys(List<String> controlKeys) {
this.controlKeys = controlKeys;
return this;
}
@Override
public String toString() {
return "CloudControlAuthRequest{" +
"userId='" + userId + '\'' +
", userCallsign='" + userCallsign + '\'' +
", controlKeys=" + controlKeys +
'}';
}
}
package com.dji.sdk.cloudapi.control;
import com.dji.sdk.common.BaseModel;
import javax.validation.constraints.NotNull;
import javax.validation.constraints.Size;
import java.util.List;
/**
* Pilot 释放云端控制权:cloud_control_release
*/
public class CloudControlReleaseRequest extends BaseModel {
@NotNull
@Size(min = 1)
private List<String> controlKeys;
public CloudControlReleaseRequest() {
}
public List<String> getControlKeys() {
return controlKeys;
}
public CloudControlReleaseRequest setControlKeys(List<String> controlKeys) {
this.controlKeys = controlKeys;
return this;
}
@Override
public String toString() {
return "CloudControlReleaseRequest{" +
"controlKeys=" + controlKeys +
'}';
}
}
......@@ -46,15 +46,15 @@ public enum ControlErrorCodeEnum implements IServicesErrorCode, IEventsErrorCode
WRONG_LENS_TYPE(327015, "Invalid camera lens type."),
DRC_ABNORMAL(514300, "DRC abnormal."),
DRC_ABNORMAL(514300, "Gateway error."),
DRC_HEARTBEAT_TIMED_OUT(514301, "DRC heartbeat timed out."),
DRC_HEARTBEAT_TIMED_OUT(514301, "Request timed out. Disconnected."),
DRC_CERTIFICATE_ABNORMAL(514302, "DRC certificate is abnormal."),
DRC_CERTIFICATE_ABNORMAL(514302, "Network certificate error. Connection failed."),
DRC_LINK_LOST(514303, "DRC link is lost."),
DRC_LINK_LOST(514303, "Network error. Disconnected."),
DRC_LINK_REFUSED(514304, "DRC link is refused."),
DRC_LINK_REFUSED(514304, "Request rejected. Connection failed."),
UNKNOWN(-1, "UNKNOWN"),
......
......@@ -11,6 +11,10 @@ public enum ControlMethodEnum {
PAYLOAD_AUTHORITY_GRAB("payload_authority_grab"),
CLOUD_CONTROL_AUTH_REQUEST("cloud_control_auth_request"),
CLOUD_CONTROL_RELEASE("cloud_control_release"),
DRC_MODE_ENTER("drc_mode_enter"),
DRC_MODE_EXIT("drc_mode_exit"),
......
......@@ -15,15 +15,15 @@ public enum DrcStatusErrorEnum implements IErrorInfo {
SUCCESS(0, "success"),
MQTT_ERR(514300, "The mqtt connection error."),
MQTT_ERR(514300, "Gateway error."),
HEARTBEAT_TIMEOUT(514301, "The heartbeat times out and the dock disconnects."),
HEARTBEAT_TIMEOUT(514301, "Request timed out. Disconnected."),
MQTT_CERTIFICATE_ERR(514302, "The mqtt certificate is abnormal and the connection fails."),
MQTT_CERTIFICATE_ERR(514302, "Network certificate error. Connection failed."),
MQTT_LOST(514303, "The dock network is abnormal and the mqtt connection is lost."),
MQTT_LOST(514303, "Network error. Disconnected."),
MQTT_REFUSE(514304, "The dock connection to mqtt server was refused.");
MQTT_REFUSE(514304, "Request rejected. Connection failed.");
private final String msg;
......
......@@ -90,7 +90,7 @@ public abstract class AbstractControlService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flightAuthorityGrab(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -103,7 +103,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> payloadAuthorityGrab(GatewayManager gateway, PayloadAuthorityGrabRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -112,12 +112,37 @@ public abstract class AbstractControlService {
}
/**
* Pilot:请求遥控器授权云端控制(遥控器弹窗确认)
* 超时放宽:飞手需在 Pilot 上点同意,默认 3s 不够。
*/
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cloudControlAuthRequest(GatewayManager gateway, CloudControlAuthRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
ControlMethodEnum.CLOUD_CONTROL_AUTH_REQUEST.getMethod(),
request,
1,
30_000L);
}
/**
* Pilot:释放云端控制授权
*/
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cloudControlRelease(GatewayManager gateway, CloudControlReleaseRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
ControlMethodEnum.CLOUD_CONTROL_RELEASE.getMethod(),
request);
}
/**
* Enter the live flight controls mode
* @param gateway
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> drcModeEnter(GatewayManager gateway, DrcModeEnterRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -130,7 +155,7 @@ public abstract class AbstractControlService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> drcModeExit(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -143,7 +168,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> takeoffToPoint(GatewayManager gateway, TakeoffToPointRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -157,7 +182,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flyToPoint(GatewayManager gateway, FlyToPointRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -184,7 +209,7 @@ public abstract class AbstractControlService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flyToPointStop(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -197,7 +222,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraModeSwitch(GatewayManager gateway, CameraModeSwitchRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -211,7 +236,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraPhotoTake(GatewayManager gateway, CameraPhotoTakeRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -252,7 +277,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraRecordingStart(GatewayManager gateway, CameraRecordingStartRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -266,7 +291,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraRecordingStop(GatewayManager gateway, CameraRecordingStopRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -280,7 +305,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraAim(GatewayManager gateway, CameraAimRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -294,7 +319,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraScreenDrag(GatewayManager gateway, CameraScreenDragRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -308,7 +333,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraFrameZoom(GatewayManager gateway, CameraFrameZoomRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -322,7 +347,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> cameraFocalLengthSet(GatewayManager gateway, CameraFocalLengthSetRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -336,7 +361,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> gimbalReset(GatewayManager gateway, GimbalResetRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -352,7 +377,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0, exclude = GatewayTypeEnum.RC)
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
public TopicServicesResponse<ServicesReplyData> cameraLookAt(GatewayManager gateway, CameraLookAtRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -366,7 +391,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0, exclude = GatewayTypeEnum.RC)
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
public TopicServicesResponse<ServicesReplyData> cameraScreenSplit(GatewayManager gateway, CameraScreenSplitRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -380,7 +405,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0, exclude = GatewayTypeEnum.RC)
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
public TopicServicesResponse<ServicesReplyData> photoStorageSet(GatewayManager gateway, PhotoStorageSetRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -394,7 +419,7 @@ public abstract class AbstractControlService {
* @param request data
* @return services_reply
*/
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0, exclude = GatewayTypeEnum.RC)
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
public TopicServicesResponse<ServicesReplyData> videoStorageSet(GatewayManager gateway, VideoStorageSetRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -590,7 +615,7 @@ public abstract class AbstractControlService {
* @param gateway
* @param request data
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
protected void droneControlDown(GatewayManager gateway, DroneControlRequest request) {
drcDownPublish.publish(
gateway.getGatewaySn(),
......@@ -613,7 +638,7 @@ public abstract class AbstractControlService {
* DRC-drone emergency stop
* @param gateway
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public void droneEmergencyStopDown(GatewayManager gateway) {
drcDownPublish.publish(
gateway.getGatewaySn(),
......@@ -637,7 +662,7 @@ public abstract class AbstractControlService {
* @param gateway
* @param request data
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public void heartBeatDown(GatewayManager gateway, HeartBeatRequest request) {
drcDownPublish.publish(
gateway.getGatewaySn(),
......
......@@ -168,7 +168,7 @@ public class AbstractDeviceService {
*/
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_WPMZ_VERSION)
public void dockWpmzVersionUpdate(TopicStateRequest<DockDroneWpmzVersion> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockWpmzVersionUpdate not implemented");
// 默认 no-op;业务实现见 sample SDKDeviceService
}
/**
......@@ -178,7 +178,7 @@ public class AbstractDeviceService {
*/
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_THERMAL_SUPPORTED_PALETTE_STYLE)
public void dockThermalSupportedPaletteStyle(TopicStateRequest<DockDroneThermalSupportedPaletteStyle> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockThermalSupportedPaletteStyle not implemented");
// optional property; ignore if not consumed by business
}
/**
......@@ -190,7 +190,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_RTH_MODE, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockDroneRthMode(TopicStateRequest<DockDroneRthMode> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockRthMode not implemented");
return ackState();
}
/**
......@@ -201,7 +201,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_CURRENT_RTH_MODE, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockDroneCurrentRthMode(TopicStateRequest<DockDroneCurrentRthMode> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockCurrentRthMode not implemented");
return ackState();
}
/**
......@@ -212,7 +212,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_COMMANDER_MODE_LOST_ACTION, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockDroneCommanderModeLostAction(TopicStateRequest<DockDroneCommanderModeLostAction> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockDroneCommanderModeLostAction not implemented");
return ackState();
}
/**
......@@ -223,7 +223,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_CURRENT_COMMANDER_FLIGHT_MODE, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockDroneCurrentCommanderFlightMode(TopicStateRequest<DockDroneCurrentCommanderFlightMode> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockDroneCurrentCommanderFlightMode not implemented");
return ackState();
}
/**
......@@ -235,7 +235,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_COMMANDER_FLIGHT_HEIGHT, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockDroneCommanderFlightHeight(TopicStateRequest<DockDroneCommanderFlightHeight> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockDroneCommanderFlightHeight not implemented");
return ackState();
}
/**
......@@ -246,7 +246,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_DRONE_MODE_CODE_REASON, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockDroneModeCodeReason(TopicStateRequest<DockDroneModeCodeReason> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockDroneModeCodeReason not implemented");
return ackState();
}
/**
......@@ -257,7 +257,7 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_1, include = GatewayTypeEnum.DOCK2)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_AND_DRONE_DONGLE_INFOS, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dongleInfos(TopicStateRequest<DongleInfos> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dongleInfos not implemented");
return ackState();
}
/**
......@@ -268,7 +268,11 @@ public class AbstractDeviceService {
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_2, include = GatewayTypeEnum.DOCK)
@ServiceActivator(inputChannel = ChannelName.INBOUND_STATE_DOCK_SILENT_MODE, outputChannel = ChannelName.OUTBOUND_STATE)
public TopicStateResponse<MqttReply> dockSilentMode(TopicStateRequest<DockSilentMode> request, MessageHeaders headers) {
throw new UnsupportedOperationException("dockSilentMode not implemented");
return ackState();
}
private static TopicStateResponse<MqttReply> ackState() {
return new TopicStateResponse<MqttReply>().setData(MqttReply.success());
}
}
package com.dji.sdk.cloudapi.livestream;
import com.dji.sdk.exception.CloudSDKException;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonValue;
......@@ -23,7 +22,10 @@ public enum VideoTypeEnum {
POINT_CLOUD("point_cloud"),
IR("ir");
IR("ir"),
/** 停流或固件未填 video_type 时常见为空串 */
UNKNOWN("");
private final String type;
......@@ -39,6 +41,6 @@ public enum VideoTypeEnum {
@JsonCreator
public static VideoTypeEnum find(String videoType) {
return Arrays.stream(values()).filter(typeEnum -> typeEnum.type.equals(videoType)).findAny()
.orElseThrow(() -> new CloudSDKException(VideoTypeEnum.class , videoType));
.orElse(UNKNOWN);
}
}
package com.dji.sdk.cloudapi.wayline;
import com.dji.sdk.exception.CloudSDKException;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonValue;
......@@ -17,7 +16,8 @@ public enum LastPointTypeEnum {
NOT_OVER_THE_HOME_POINT(1),
;
/** 固件未就绪/无效时常见哨兵值(如 65535) */
UNKNOWN(65535);
private final int type;
......@@ -33,7 +33,7 @@ public enum LastPointTypeEnum {
@JsonCreator
public static LastPointTypeEnum find(int type) {
return Arrays.stream(values()).filter(typeEnum -> typeEnum.type == type).findAny()
.orElseThrow(() -> new CloudSDKException(LastPointTypeEnum.class, type));
.orElse(UNKNOWN);
}
}
......@@ -43,7 +43,7 @@ public abstract class AbstractWaylineService {
*/
@ServiceActivator(inputChannel = ChannelName.INBOUND_EVENTS_DEVICE_EXIT_HOMING_NOTIFY, outputChannel = ChannelName.OUTBOUND_EVENTS)
public TopicEventsResponse<MqttReply> deviceExitHomingNotify(TopicEventsRequest<DeviceExitHomingNotify> request, MessageHeaders headers) {
throw new UnsupportedOperationException("deviceExitHomingNotify not implemented");
return new TopicEventsResponse<MqttReply>().setData(MqttReply.success());
}
/**
......@@ -73,7 +73,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(deprecated = CloudSDKVersionEnum.V0_0_1, exclude = GatewayTypeEnum.RC)
@CloudSDKVersion(deprecated = CloudSDKVersionEnum.V0_0_1)
public TopicServicesResponse<ServicesReplyData> flighttaskCreate(GatewayManager gateway, FlighttaskCreateRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -86,7 +86,8 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
/** Pilot RC 可作为网关接收航线任务(模拟机场);设备固件需实测是否响应 */
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flighttaskPrepare(GatewayManager gateway, FlighttaskPrepareRequest request) {
validPrepareParam(request);
return servicesPublish.publish(
......@@ -101,7 +102,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flighttaskExecute(GatewayManager gateway, FlighttaskExecuteRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -115,7 +116,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flighttaskUndo(GatewayManager gateway, FlighttaskUndoRequest request) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -128,7 +129,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flighttaskPause(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -140,7 +141,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> flighttaskRecovery(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -152,7 +153,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> returnHome(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -164,7 +165,7 @@ public abstract class AbstractWaylineService {
* @param gateway
* @return services_reply
*/
@CloudSDKVersion(exclude = GatewayTypeEnum.RC)
@CloudSDKVersion
public TopicServicesResponse<ServicesReplyData> returnHomeCancel(GatewayManager gateway) {
return servicesPublish.publish(
gateway.getGatewaySn(),
......@@ -190,8 +191,8 @@ public abstract class AbstractWaylineService {
*/
@ServiceActivator(inputChannel = ChannelName.INBOUND_EVENTS_RETURN_HOME_INFO, outputChannel = ChannelName.OUTBOUND_EVENTS)
@CloudSDKVersion(since = CloudSDKVersionEnum.V1_0_0)
public TopicRequestsResponse<MqttReply> returnHomeInfo(TopicRequestsRequest<ReturnHomeInfo> request, MessageHeaders headers) {
throw new UnsupportedOperationException("returnHomeInfo not implemented");
public TopicEventsResponse<MqttReply> returnHomeInfo(TopicEventsRequest<ReturnHomeInfo> request, MessageHeaders headers) {
return new TopicEventsResponse<MqttReply>().setData(MqttReply.success());
}
private void validPrepareParam(FlighttaskPrepareRequest request) {
......
......@@ -10,6 +10,7 @@ import com.dji.sdk.exception.CloudSDKErrorEnum;
import com.dji.sdk.exception.CloudSDKException;
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
/**
......@@ -24,6 +25,14 @@ public class SDKManager {
private static final ConcurrentHashMap<String, GatewayManager> SDK_MAP = new ConcurrentHashMap<>(16);
public static boolean isRegistered(String gatewaySn) {
return gatewaySn != null && SDK_MAP.containsKey(gatewaySn);
}
public static Optional<GatewayManager> findDeviceSDK(String gatewaySn) {
return Optional.ofNullable(SDK_MAP.get(gatewaySn));
}
public static GatewayManager getDeviceSDK(String gatewaySn) {
if (SDK_MAP.containsKey(gatewaySn)) {
return SDK_MAP.get(gatewaySn);
......
package com.dji.sdk.config.version;
import com.dji.sdk.exception.CloudSDKVersionException;
import com.fasterxml.jackson.annotation.JsonValue;
import java.util.Arrays;
......@@ -71,6 +70,8 @@ public enum Dock2ThingVersionEnum implements IThingVersion {
public static Dock2ThingVersionEnum find(String thingVersion) {
return Arrays.stream(values()).filter(thingVersionEnum -> thingVersionEnum.thingVersion.equals(thingVersion))
.findAny().orElseThrow(() -> new CloudSDKVersionException(thingVersion));
.findAny()
// null/未知物模型版本:回落最新,避免重启临时注册或 Redis 缺版本时直接炸
.orElse(values()[values().length - 1]);
}
}
package com.dji.sdk.config.version;
import com.dji.sdk.exception.CloudSDKVersionException;
import com.fasterxml.jackson.annotation.JsonValue;
import java.util.Arrays;
......@@ -64,6 +63,7 @@ public enum Dock3ThingVersionEnum implements IThingVersion {
public static Dock3ThingVersionEnum find(String thingVersion) {
return Arrays.stream(values()).filter(thingVersionEnum -> thingVersionEnum.thingVersion.equals(thingVersion))
.findAny().orElseThrow(() -> new CloudSDKVersionException(thingVersion));
.findAny()
.orElse(values()[values().length - 1]);
}
}
package com.dji.sdk.config.version;
import com.dji.sdk.exception.CloudSDKVersionException;
import com.fasterxml.jackson.annotation.JsonValue;
import java.util.Arrays;
......@@ -74,6 +73,7 @@ public enum DockThingVersionEnum implements IThingVersion {
public static DockThingVersionEnum find(String thingVersion) {
return Arrays.stream(values()).filter(thingVersionEnum -> thingVersionEnum.thingVersion.equals(thingVersion))
.findAny().orElseThrow(() -> new CloudSDKVersionException(thingVersion));
.findAny()
.orElse(values()[values().length - 1]);
}
}
package com.dji.sdk.config.version;
import com.dji.sdk.exception.CloudSDKVersionException;
import com.fasterxml.jackson.annotation.JsonValue;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Arrays;
import java.util.Optional;
......@@ -65,8 +62,6 @@ public enum DroneThingVersionEnum implements IThingVersion {
V1_2_3("1.2.3", CloudSDKVersionEnum.V1_0_3),
;
private static final Logger log = LoggerFactory.getLogger(DroneThingVersionEnum.class);
private final String thingVersion;
private final CloudSDKVersionEnum cloudSDKVersion;
......@@ -88,9 +83,6 @@ public enum DroneThingVersionEnum implements IThingVersion {
public static DroneThingVersionEnum find(String thingVersion) {
Optional<DroneThingVersionEnum> opt = Arrays.stream(values())
.filter(thingVersionEnum -> thingVersionEnum.thingVersion.equals(thingVersion)).findAny();
if (opt.isPresent()) {
return opt.get();
}
throw new CloudSDKVersionException(thingVersion);
return opt.orElse(values()[values().length - 1]);
}
}
......@@ -29,10 +29,7 @@ public class GatewayManager {
public GatewayManager(String gatewaySn, String droneSn, GatewayTypeEnum gatewayType, String gatewayThingVersion, String droneThingVersion) {
this(gatewaySn, droneSn, gatewayType);
this.gatewayThingVersion = new GatewayThingVersion(gatewayType, gatewayThingVersion);
if (GatewayTypeEnum.RC == gatewayType) {
this.sdkVersion = CloudSDKVersionEnum.V0_0_1;
return;
}
// Pilot RC 作为虚拟机场时,按物模型版本计算 CloudSDK 版本,避免被钉死在 V0_0_1 而无法下发航线/飞控
if (Objects.isNull(droneThingVersion)) {
this.sdkVersion = this.gatewayThingVersion.getCloudSDKVersion();
return;
......
......@@ -2,11 +2,14 @@ package com.dji.sdk.mqtt.osd;
import com.dji.sdk.cloudapi.device.PayloadModelConst;
import com.dji.sdk.common.Common;
import com.dji.sdk.config.version.GatewayManager;
import com.dji.sdk.common.SDKManager;
import com.dji.sdk.config.version.GatewayManager;
import com.dji.sdk.config.version.GatewayTypeEnum;
import com.dji.sdk.exception.CloudSDKException;
import com.dji.sdk.mqtt.ChannelName;
import com.fasterxml.jackson.core.type.TypeReference;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.dsl.IntegrationFlow;
......@@ -31,6 +34,8 @@ import static com.dji.sdk.mqtt.TopicConst.*;
@Configuration
public class OsdRouter {
private static final Logger log = LoggerFactory.getLogger(OsdRouter.class);
@Bean
public IntegrationFlow osdRouterFlow() {
return IntegrationFlows
......@@ -45,9 +50,10 @@ public class OsdRouter {
}
}, null)
.<TopicOsdRequest>handle((response, headers) -> {
GatewayManager gateway = SDKManager.getDeviceSDK(response.getGateway());
OsdDeviceTypeEnum typeEnum = OsdDeviceTypeEnum.find(gateway.getType(), response.getFrom().equals(response.getGateway()));
Map<String, Object> data = (Map<String, Object>) response.getData();
boolean isGateway = response.getFrom().equals(response.getGateway());
GatewayTypeEnum gatewayType = resolveGatewayType(response.getGateway(), isGateway, data);
OsdDeviceTypeEnum typeEnum = OsdDeviceTypeEnum.find(gatewayType, isGateway);
if (!typeEnum.isGateway()) {
List payloadData = (List) data.getOrDefault(PayloadModelConst.PAYLOAD_KEY, new ArrayList<>());
PayloadModelConst.getAllIndexWithPosition().stream().filter(data::containsKey)
......@@ -61,4 +67,34 @@ public class OsdRouter {
.get();
}
}
\ No newline at end of file
/**
* 优先用已注册网关类型;未注册时(重启后 status 尚未回来)仅按 OSD 字段推断,
* 不调用 registerDevice(null thingVersion),避免 CloudSDKVersionException。
* 正式注册交给 osd* 里的 recoverDeviceOnline / update_topo。
*/
private GatewayTypeEnum resolveGatewayType(String gatewaySn, boolean isGateway, Map<String, Object> data) {
return SDKManager.findDeviceSDK(gatewaySn)
.map(GatewayManager::getType)
.orElseGet(() -> {
GatewayTypeEnum inferred = inferGatewayType(data, isGateway);
log.debug("OSD from unregistered gateway {}, infer type {}", gatewaySn, inferred);
return inferred;
});
}
private static GatewayTypeEnum inferGatewayType(Map<String, Object> data, boolean isGatewayPayload) {
boolean looksLikeDock = data.containsKey("drone_in_dock")
|| data.containsKey("cover_state")
|| data.containsKey("flighttask_step_code")
|| data.containsKey("drone_charge_state")
|| data.containsKey("air_conditioner")
|| data.containsKey("job_number")
|| data.containsKey("backup_battery")
|| data.containsKey("alternate_land_point");
if (looksLikeDock || !isGatewayPayload) {
return GatewayTypeEnum.DOCK2;
}
return GatewayTypeEnum.RC;
}
}
......@@ -6,7 +6,8 @@ import com.dji.sdk.cloudapi.property.DockDroneCommanderFlightHeight;
import com.dji.sdk.cloudapi.property.DockDroneCommanderModeLostAction;
import com.dji.sdk.cloudapi.property.DockDroneOfflineMapEnable;
import com.dji.sdk.cloudapi.property.DockDroneRthMode;
import com.dji.sdk.exception.CloudSDKException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Arrays;
import java.util.Collections;
......@@ -50,8 +51,16 @@ public enum DockStateDataKeyEnum {
SILENT_MODE(Set.of("silent_mode"), DockSilentMode.class),
/** 图传拓扑(多机场等);暂无业务消费,走 DEFAULT */
WIRELESS_LINK_TOPO(Set.of("wireless_link_topo", "wireless_link_topo_all"), Object.class),
/** 固件新增未收录字段:不抛异常,路由到 DEFAULT */
UNKNOWN(Collections.emptySet(), Object.class),
;
private static final Logger log = LoggerFactory.getLogger(DockStateDataKeyEnum.class);
private final Set<String> keys;
private final Class classType;
......@@ -71,8 +80,13 @@ public enum DockStateDataKeyEnum {
}
public static DockStateDataKeyEnum find(Set<String> keys) {
return Arrays.stream(values()).filter(keyEnum -> !Collections.disjoint(keys, keyEnum.keys)).findAny()
.orElseThrow(() -> new CloudSDKException(DockStateDataKeyEnum.class, keys));
return Arrays.stream(values())
.filter(keyEnum -> keyEnum != UNKNOWN && !Collections.disjoint(keys, keyEnum.keys))
.findAny()
.orElseGet(() -> {
log.debug("Dock state unknown keys {}, route to DEFAULT", keys);
return UNKNOWN;
});
}
}
......@@ -2,7 +2,8 @@ package com.dji.sdk.mqtt.state;
import com.dji.sdk.cloudapi.device.*;
import com.dji.sdk.cloudapi.livestream.RcLivestreamAbilityUpdate;
import com.dji.sdk.exception.CloudSDKException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Arrays;
import java.util.Collections;
......@@ -25,8 +26,28 @@ public enum RcStateDataKeyEnum {
LIVE_STATUS(Set.of("live_status"), RcLiveStatus.class),
PAYLOAD_FIRMWARE(PayloadModelConst.getAllModelWithPosition(), PayloadFirmwareVersion.class),
/** RC / 飞行器也会上报 4G Dongle;与机场共用 DongleInfos 模型 */
DONGLE_INFOS(Set.of("dongle_infos"), DongleInfos.class),
/** Pilot/机场共用:当前指点飞模式 */
CURRENT_COMMANDER_FLIGHT_MODE(Set.of("current_commander_flight_mode"), Object.class),
/**
* Pilot 云端控制授权状态(文档亦见 is_cloud_control_auth)。
* 暂无强类型消费,走 DEFAULT,避免刷 ERROR。
*/
CLOUD_CONTROL_AUTH(Set.of("cloud_control_auth", "is_cloud_control_auth"), Object.class),
/** 图传拓扑等固件新增字段;暂无业务消费,走 DEFAULT */
WIRELESS_LINK_TOPO(Set.of("wireless_link_topo", "wireless_link_topo_all"), Object.class),
/** 固件新增未收录字段:不抛异常,路由到 DEFAULT */
UNKNOWN(Collections.emptySet(), Object.class),
;
private static final Logger log = LoggerFactory.getLogger(RcStateDataKeyEnum.class);
private final Set<String> keys;
private final Class classType;
......@@ -46,8 +67,13 @@ public enum RcStateDataKeyEnum {
}
public static RcStateDataKeyEnum find(Set<String> keys) {
return Arrays.stream(values()).filter(keyEnum -> !Collections.disjoint(keys, keyEnum.keys)).findAny()
.orElseThrow(() -> new CloudSDKException(RcStateDataKeyEnum.class, keys));
return Arrays.stream(values())
.filter(keyEnum -> keyEnum != UNKNOWN && !Collections.disjoint(keys, keyEnum.keys))
.findAny()
.orElseGet(() -> {
log.debug("RC state unknown keys {}, route to DEFAULT", keys);
return UNKNOWN;
});
}
}
......@@ -77,15 +77,20 @@ public class StateRouter {
private Class getTypeReference(String gatewaySn, Object data) {
Set<String> keys = ((Map<String, Object>) data).keySet();
switch (SDKManager.getDeviceSDK(gatewaySn).getType()) {
case RC:
return RcStateDataKeyEnum.find(keys).getClassType();
case DOCK:
case DOCK2:
case DOCK3:
return DockStateDataKeyEnum.find(keys).getClassType();
default:
throw new CloudSDKException(CloudSDKErrorEnum.WRONG_DATA, "Unexpected value: " + SDKManager.getDeviceSDK(gatewaySn).getType());
try {
switch (SDKManager.getDeviceSDK(gatewaySn).getType()) {
case RC:
return RcStateDataKeyEnum.find(keys).getClassType();
case DOCK:
case DOCK2:
case DOCK3:
return DockStateDataKeyEnum.find(keys).getClassType();
default:
throw new CloudSDKException(CloudSDKErrorEnum.WRONG_DATA, "Unexpected value: " + SDKManager.getDeviceSDK(gatewaySn).getType());
}
} catch (CloudSDKException e) {
// 兼容旧包 / 未收录字段:降级为 Object,由 StateDataKeyEnum.UNKNOWN → DEFAULT 消化
return Object.class;
}
}
}
\ No newline at end of file
......@@ -7,6 +7,7 @@ import com.dji.sample.manage.service.IDeviceRedisService;
import com.dji.sample.manage.service.IDeviceService;
import com.dji.sdk.cloudapi.device.DeviceDomainEnum;
import com.dji.sdk.common.SDKManager;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;
......@@ -18,6 +19,7 @@ import java.util.Optional;
* @date 2021/11/24
* @version 0.1
*/
@Slf4j
@Component
public class ApplicationBootInitial implements CommandLineRunner {
......@@ -41,12 +43,19 @@ public class ApplicationBootInitial implements CommandLineRunner {
.stream()
.map(key -> key.substring(start))
.map(deviceRedisService::getDeviceOnline)
.filter(Optional::isPresent)
.map(Optional::get)
.filter(device -> DeviceDomainEnum.DRONE != device.getDomain())
.forEach(device -> deviceService.subDeviceOnlineSubscribeTopic(
SDKManager.registerDevice(device.getDeviceSn(), device.getChildDeviceSn(), device.getDomain(),
device.getType(), device.getSubType(), device.getThingVersion(),
deviceRedisService.getDeviceOnline(device.getChildDeviceSn()).map(DeviceDTO::getThingVersion).orElse(null))));
.forEach(device -> {
try {
deviceService.subDeviceOnlineSubscribeTopic(
SDKManager.registerDevice(device.getDeviceSn(), device.getChildDeviceSn(), device.getDomain(),
device.getType(), device.getSubType(), device.getThingVersion(),
deviceRedisService.getDeviceOnline(device.getChildDeviceSn()).map(DeviceDTO::getThingVersion).orElse(null)));
} catch (Exception e) {
log.error("Boot re-subscribe failed for {}: {}", device.getDeviceSn(), e.getMessage());
}
});
}
}
\ No newline at end of file
}
......@@ -5,13 +5,12 @@ import com.dji.sample.component.redis.RedisOpsUtils;
import com.dji.sample.manage.model.dto.DeviceDTO;
import com.dji.sample.manage.service.IDeviceRedisService;
import com.dji.sample.manage.service.IDeviceService;
import com.dji.sdk.common.SDKManager;
import com.dji.sdk.cloudapi.device.DeviceDomainEnum;
import com.dji.sdk.common.SDKManager;
import com.dji.sdk.mqtt.IMqttTopicService;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
......@@ -41,8 +40,6 @@ public class GlobalScheduleService {
@Autowired
private MqttPahoMessageDrivenChannelAdapter mqttInbound;
@Autowired
private ObjectMapper mapper;
/**
* Check the status of the devices every 30 seconds. It is recommended to use cache.
*/
......@@ -70,54 +67,25 @@ public class GlobalScheduleService {
}
/**
* Check MQTT connection/subscriptions every 30 seconds and recover missing dynamic subscriptions.
* Recover MQTT adapter if down; always re-sync online gateways into SDKManager.
* Wildcard +/osd can deliver traffic before status re-login after restart.
*/
@Scheduled(initialDelay = 30, fixedRate = 30, timeUnit = TimeUnit.SECONDS)
@Scheduled(initialDelay = 5, fixedRate = 30, timeUnit = TimeUnit.SECONDS)
private void mqttConnectionCheck() {
if (!mqttInbound.isRunning()) {
log.warn("MQTT adapter is not running, attempting to restart and re-subscribe...");
try {
mqttInbound.start();
resubscribeOnlineDevices();
} catch (Exception e) {
log.error("Failed to restart MQTT adapter", e);
return;
}
return;
}
String[] currentTopics = topicService.getSubscribedTopic();
int start = RedisConst.DEVICE_ONLINE_PREFIX.length();
boolean missingSubscription = false;
for (String key : RedisOpsUtils.getAllKeys(RedisConst.DEVICE_ONLINE_PREFIX + "*")) {
String sn = key.substring(start);
String osdTopic = "thing/product/" + sn + "/osd";
if (!containsTopic(currentTopics, osdTopic) && !containsTopic(currentTopics, "thing/product/+/osd")) {
missingSubscription = true;
break;
}
}
if (missingSubscription) {
log.warn("Detected missing MQTT subscriptions after reconnect, re-subscribing online devices...");
resubscribeOnlineDevices();
}
}
private boolean containsTopic(String[] topics, String expected) {
if (topics == null || topics.length == 0 || !StringUtils.hasText(expected)) {
return false;
}
for (String topic : topics) {
if (expected.equals(topic)) {
return true;
}
}
return false;
resubscribeOnlineDevices();
}
/**
* Re-subscribe all currently online non-drone devices.
* Re-subscribe all currently online non-drone devices and ensure SDKManager registration.
*/
private void resubscribeOnlineDevices() {
int start = RedisConst.DEVICE_ONLINE_PREFIX.length();
......@@ -131,10 +99,10 @@ public class GlobalScheduleService {
try {
if (DeviceDomainEnum.DRONE == device.getDomain()) {
deviceService.subDroneOnlineSubscribeTopic(device.getDeviceSn());
log.info("Re-subscribed drone device: {}", device.getDeviceSn());
return;
}
boolean wasRegistered = SDKManager.isRegistered(device.getDeviceSn());
String childSn = device.getChildDeviceSn();
String childThingVersion = StringUtils.hasText(childSn)
? deviceRedisService.getDeviceOnline(childSn).map(DeviceDTO::getThingVersion).orElse(null)
......@@ -144,11 +112,13 @@ public class GlobalScheduleService {
SDKManager.registerDevice(device.getDeviceSn(), childSn,
device.getDomain(), device.getType(), device.getSubType(),
device.getThingVersion(), childThingVersion));
log.info("Re-subscribed device: {}", device.getDeviceSn());
if (!wasRegistered) {
log.info("Re-registered online gateway into SDKManager: {}", device.getDeviceSn());
}
} catch (Exception e) {
log.error("Failed to re-subscribe device: {}", device.getDeviceSn(), e);
}
});
}
}
\ No newline at end of file
}
......@@ -7,6 +7,7 @@ import com.dji.sample.component.mqtt.model.MqttProtocolEnum;
import com.dji.sample.component.mqtt.model.MqttUseEnum;
import com.dji.sdk.cloudapi.control.DrcModeMqttBroker;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
......@@ -23,6 +24,7 @@ import java.util.Map;
* @date 2021/11/10
* @version 0.1
*/
@Slf4j
@Configuration
@Data
@ConfigurationProperties
......@@ -91,28 +93,71 @@ public class MqttPropertyConfiguration {
}
/**
* Get the connection parameters of the mqtt client of the drc link.
* @param clientId
* @param username
* @param age The validity period of the token. unit: s
* @param map Custom data added in token.
* @return
* Web 前端 DRC 连接(mqtt.js):需要带 scheme 的完整 URL,如 wss://host:443/mqtt
*/
public static DrcModeMqttBroker getMqttBrokerWithDrc(String clientId, String username, Long age, Map<String, ?> map) {
return buildDrcBroker(clientId, username, age, map, false);
}
/**
* 下发给机场/Pilot 的 DRC Broker(官方格式):address=host:port,无 wss:// 前缀。
*/
public static DrcModeMqttBroker getMqttBrokerWithDrcForDevice(String clientId, String username, Long age, Map<String, ?> map) {
return buildDrcBroker(clientId, username, age, map, true);
}
private static DrcModeMqttBroker buildDrcBroker(String clientId, String username, Long age, Map<String, ?> map, boolean forDevice) {
if (!mqtt.containsKey(MqttUseEnum.DRC)) {
throw new RuntimeException("Please configure the drc link parameters of mqtt in the backend configuration file first.");
}
Algorithm algorithm = JwtUtil.algorithm;
String token = JwtUtil.createToken(map, age, algorithm, null, null);
MqttClientOptions drcOptions = mqtt.get(MqttUseEnum.DRC);
String address;
boolean enableTls;
if (forDevice) {
String host = StringUtils.hasText(drcOptions.getDeviceHost())
? drcOptions.getDeviceHost().trim()
: drcOptions.getHost().trim();
int port;
if (drcOptions.getDevicePort() != null) {
port = drcOptions.getDevicePort();
} else if (drcOptions.getProtocol() == MqttProtocolEnum.WSS
|| drcOptions.getProtocol() == MqttProtocolEnum.WS) {
// Web 用 WSS:443;设备官方示例为 MQTTS:8883。未单独配置时默认 8883,避免把 443/wss 路径塞给 Pilot。
port = 8883;
log.warn("mqtt.DRC.device-port unset; using {}:8883 for device DRC (set device-port to override)", host);
} else {
port = drcOptions.getPort();
}
address = host + ":" + port;
if (drcOptions.getDeviceEnableTls() != null) {
enableTls = drcOptions.getDeviceEnableTls();
} else if (port == 8883) {
enableTls = true;
} else if (drcOptions.getDevicePort() != null) {
// 显式配了非 8883 设备端口(如 1883/54418):默认明文,勿因 Web 的 WSS 误开 TLS
enableTls = false;
} else {
enableTls = drcOptions.getProtocol() == MqttProtocolEnum.MQTTS;
}
} else {
address = getMqttAddress(drcOptions);
enableTls = drcOptions.getProtocol() == MqttProtocolEnum.WSS
|| drcOptions.getProtocol() == MqttProtocolEnum.MQTTS;
}
log.info("Build DRC mqtt broker: forDevice={}, address={}, enableTls={}, clientId={}, username={}",
forDevice, address, enableTls, clientId, username);
return new DrcModeMqttBroker()
.setAddress(getMqttAddress(mqtt.get(MqttUseEnum.DRC)))
.setAddress(address)
.setUsername(username)
.setClientId(clientId)
.setExpireTime(System.currentTimeMillis() / 1000 + age)
.setPassword(token)
.setEnableTls(false);
.setEnableTls(enableTls);
}
......
......@@ -25,6 +25,22 @@ public class MqttClientOptions {
private String path;
/**
* 下发给机场/Pilot 的 DRC 地址主机(官方格式 host:port,无 scheme)。
* 不填则回退到 {@link #host}。
*/
private String deviceHost;
/**
* 下发给设备的 DRC 端口。官方示例为 MQTTS 8883;不填则回退到 {@link #port}。
*/
private Integer devicePort;
/**
* 设备侧是否 TLS。不填时:MQTTS/WSS 为 true,否则 false。
*/
private Boolean deviceEnableTls;
/**
* The topic to subscribe to immediately when client connects. Only required for basic link.
*/
private String inboundTopic;
......
......@@ -53,6 +53,12 @@ public final class RedisConst {
public static final String OSD_PREFIX = "osd" + DELIMITER;
/** 飞行器 WPMZ(航线任务库)版本,完整 key = wpmz_version:{sn} */
public static final String WPMZ_VERSION_PREFIX = "wpmz_version" + DELIMITER;
/** 飞行器 mode_code_reason,完整 key = mode_code_reason:{sn} */
public static final String MODE_CODE_REASON_PREFIX = "mode_code_reason" + DELIMITER;
public static final String MEDIA_FILE_PREFIX = "media_file" + DELIMITER;
public static final String MEDIA_HIGHEST_PRIORITY_PREFIX = "media_highest_priority" + DELIMITER;
......
......@@ -93,10 +93,30 @@ public class MyWebSocketHandler extends WebSocketDefaultHandler {
String payload = message.getPayload().toString();
log.debug("Received message from session {}: {}", sessionId, payload);
// 明文心跳(网关/旧客户端常见),勿按 JSON 解析
String trimmed = payload == null ? "" : payload.trim();
if ("ping".equalsIgnoreCase(trimmed)) {
if (session.isOpen()) {
session.sendMessage(new org.springframework.web.socket.TextMessage("pong"));
}
return;
}
if ("pong".equalsIgnoreCase(trimmed)) {
return;
}
try {
JsonNode jsonNode = objectMapper.readTree(payload);
String type = jsonNode.has("type") ? jsonNode.get("type").asText() : null;
if ("ping".equalsIgnoreCase(type)) {
if (session.isOpen()) {
session.sendMessage(new org.springframework.web.socket.TextMessage(
objectMapper.writeValueAsString(Map.of("type", "pong"))));
}
return;
}
if ("auth".equals(type) && !Boolean.TRUE.equals(isAuthenticated)) {
String token = jsonNode.has("token") ? jsonNode.get("token").asText() : null;
if (StringUtils.hasText(token)) {
......
......@@ -19,6 +19,21 @@ public enum BizCodeEnum {
DOCK_OSD("dock_osd"),
/** 机场直播状态 state.live_status */
DOCK_LIVE_STATUS("dock_live_status"),
/** 遥控器直播状态 state.live_status */
RC_LIVE_STATUS("rc_live_status"),
/** AirSense 有人机告警 events.airsense_warning */
AIRSENSE_WARNING("airsense_warning"),
/** 飞行器进入当前状态的原因 state.mode_code_reason */
MODE_CODE_REASON("mode_code_reason"),
/** 航线任务库 WPMZ 版本 state.wpmz_version */
WPMZ_VERSION("wpmz_version"),
MAP_ELEMENT_CREATE("map_element_create"),
MAP_ELEMENT_UPDATE("map_element_update"),
......
......@@ -75,6 +75,22 @@ public class DockController {
return result;
}
/** Pilot:申请云端控制授权(遥控器弹窗) */
@PostMapping("/{sn}/authority/cloud-control")
public HttpResultResponse requestCloudControlAuth(HttpServletRequest request, @PathVariable String sn) {
HttpResultResponse result = controlService.requestCloudControlAuth(sn);
operateRecordService.record(request, OperateRecordTypeEnum.SEIZE_FLIGHT_AUTHORITY, sn, null);
return result;
}
/** Pilot:释放云端控制授权 */
@DeleteMapping("/{sn}/authority/cloud-control")
public HttpResultResponse releaseCloudControlAuth(HttpServletRequest request, @PathVariable String sn) {
HttpResultResponse result = controlService.releaseCloudControlAuth(sn);
operateRecordService.record(request, OperateRecordTypeEnum.SEIZE_FLIGHT_AUTHORITY, sn, null);
return result;
}
@PostMapping("/{sn}/authority/payload")
public HttpResultResponse seizePayloadAuthority(HttpServletRequest request,@PathVariable String sn, @Valid @RequestBody DronePayloadParam param) {
HttpResultResponse result = controlService.seizeAuthority(sn, DroneAuthorityEnum.PAYLOAD, param);
......
......@@ -62,6 +62,16 @@ public interface IControlService {
HttpResultResponse seizeAuthority(String sn, DroneAuthorityEnum authority, DronePayloadParam param);
/**
* Pilot:向遥控器申请云端控制授权(遥控器弹窗确认)
*/
HttpResultResponse requestCloudControlAuth(String sn);
/**
* Pilot:释放云端控制授权
*/
HttpResultResponse releaseCloudControlAuth(String sn);
/**
* Control the payload of the drone.
* @param param
* @return
......
......@@ -27,6 +27,9 @@ public class CameraFocalLengthSetImpl extends PayloadCommandsHandler {
@Override
public boolean canPublish(String deviceSn) {
super.canPublish(deviceSn);
if (osdCamera == null) {
return true;
}
if (CameraStateEnum.WORKING == osdCamera.getPhotoState()) {
return false;
}
......
......@@ -24,6 +24,9 @@ public class CameraModeSwitchImpl extends PayloadCommandsHandler {
@Override
public boolean canPublish(String deviceSn) {
super.canPublish(deviceSn);
if (osdCamera == null) {
return true;
}
return param.getCameraMode() != osdCamera.getCameraMode()
&& CameraStateEnum.IDLE == osdCamera.getPhotoState()
&& CameraStateEnum.IDLE == osdCamera.getRecordingState();
......
......@@ -17,6 +17,10 @@ public class CameraPhotoTakeImpl extends PayloadCommandsHandler {
@Override
public boolean canPublish(String deviceSn) {
super.canPublish(deviceSn);
// Pilot OsdRcDrone 无 camera OSD,跳过状态门禁
if (osdCamera == null) {
return true;
}
return CameraStateEnum.WORKING != osdCamera.getPhotoState() && osdCamera.getRemainPhotoNum() > 0;
}
}
......@@ -18,6 +18,9 @@ public class CameraRecordingStartImpl extends PayloadCommandsHandler {
@Override
public boolean canPublish(String deviceSn) {
super.canPublish(deviceSn);
if (osdCamera == null) {
return true;
}
return CameraModeEnum.VIDEO == osdCamera.getCameraMode()
&& CameraStateEnum.IDLE == osdCamera.getRecordingState()
&& osdCamera.getRemainRecordDuration() > 0;
......
......@@ -17,6 +17,9 @@ public class CameraRecordingStopImpl extends PayloadCommandsHandler {
@Override
public boolean canPublish(String deviceSn) {
super.canPublish(deviceSn);
if (osdCamera == null) {
return true;
}
return CameraStateEnum.WORKING == osdCamera.getRecordingState();
}
}
......@@ -13,6 +13,8 @@ import com.dji.sample.manage.model.dto.DeviceDTO;
import com.dji.sample.manage.service.IDevicePayloadService;
import com.dji.sample.manage.service.IDeviceRedisService;
import com.dji.sample.manage.service.IDeviceService;
import com.dji.sdk.cloudapi.control.CloudControlAuthRequest;
import com.dji.sdk.cloudapi.control.CloudControlReleaseRequest;
import com.dji.sdk.cloudapi.control.ControlMethodEnum;
import com.dji.sdk.cloudapi.control.FlyToPointRequest;
import com.dji.sdk.cloudapi.control.PayloadAuthorityGrabRequest;
......@@ -20,6 +22,7 @@ import com.dji.sdk.cloudapi.control.TakeoffToPointRequest;
import com.dji.sdk.cloudapi.control.api.AbstractControlService;
import com.dji.sdk.cloudapi.debug.DebugMethodEnum;
import com.dji.sdk.cloudapi.debug.api.AbstractDebugService;
import com.dji.sdk.cloudapi.device.DeviceDomainEnum;
import com.dji.sdk.cloudapi.device.DockModeCodeEnum;
import com.dji.sdk.cloudapi.device.DroneModeCodeEnum;
import com.dji.sdk.cloudapi.device.PayloadIndex;
......@@ -27,6 +30,7 @@ import com.dji.sdk.cloudapi.wayline.api.AbstractWaylineService;
import com.dji.sdk.common.HttpResultResponse;
import com.dji.sdk.common.SDKManager;
import com.dji.sdk.exception.CloudSDKErrorEnum;
import com.dji.sdk.exception.CloudSDKException;
import com.dji.sdk.mqtt.services.ServicesReplyData;
import com.dji.sdk.mqtt.services.TopicServicesResponse;
import com.fasterxml.jackson.databind.ObjectMapper;
......@@ -35,10 +39,12 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Collections;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
......@@ -176,7 +182,17 @@ public class ControlServiceImpl implements IControlService {
return;
}
Optional<DeviceDTO> dockOpt = deviceRedisService.getDeviceOnline(dockSn);
if (dockOpt.isEmpty() || DockModeCodeEnum.IDLE != deviceService.getDockMode(dockSn)) {
if (dockOpt.isEmpty()) {
throw new RuntimeException("The current state does not support takeoff.");
}
DeviceDTO gateway = dockOpt.get();
// Pilot RC 模拟机场:无 Dock mode_code,仅要求网关在线且子机存在
if (DeviceDomainEnum.REMOTER_CONTROL == gateway.getDomain()) {
if (!StringUtils.hasText(gateway.getChildDeviceSn())
|| !deviceRedisService.checkDeviceOnline(gateway.getChildDeviceSn())) {
throw new RuntimeException("The current state does not support takeoff.");
}
} else if (DockModeCodeEnum.IDLE != deviceService.getDockMode(dockSn)) {
throw new RuntimeException("The current state does not support takeoff.");
}
......@@ -212,8 +228,13 @@ public class ControlServiceImpl implements IControlService {
RedisOpsUtils.setWithExpire("templateTask:"+param.getFlightId(), takeoffToPointRequest, 24*60*60);
TopicServicesResponse<ServicesReplyData> response = abstractControlService.takeoffToPoint(
SDKManager.getDeviceSDK(sn), takeoffToPointRequest);
TopicServicesResponse<ServicesReplyData> response;
try {
response = abstractControlService.takeoffToPoint(
SDKManager.getDeviceSDK(sn), takeoffToPointRequest);
} catch (CloudSDKException e) {
throw new RuntimeException(mqttNoReplyHint(sn, "takeoff_to_point"), e);
}
ServicesReplyData reply = response.getData();
return reply.getResult().isSuccess() ?
HttpResultResponse.success()
......@@ -225,6 +246,11 @@ public class ControlServiceImpl implements IControlService {
TopicServicesResponse<ServicesReplyData> response;
switch (authority) {
case FLIGHT:
// Pilot RC:topo 上的 controlSource=A 表示遥控器本控,不是云端已授权。
// 不能走 Dock 的 checkAuthorityFlight,必须走 cloud_control_auth_request。
if (isRemoterControlGateway(sn)) {
return requestCloudControlAuth(sn);
}
if (deviceService.checkAuthorityFlight(sn)) {
return HttpResultResponse.success();
}
......@@ -247,6 +273,55 @@ public class ControlServiceImpl implements IControlService {
: HttpResultResponse.error(serviceReply.getResult());
}
private boolean isRemoterControlGateway(String sn) {
return deviceRedisService.getDeviceOnline(sn)
.map(d -> DeviceDomainEnum.REMOTER_CONTROL == d.getDomain())
.orElse(false);
}
private String mqttNoReplyHint(String gatewaySn, String method) {
if (isRemoterControlGateway(gatewaySn)) {
return "Pilot 遥控器未应答 " + method + "(211001)。"
+ "请先在遥控器确认云端控制授权弹窗;确认 dock_sn 为遥控器 SN,"
+ "并用 MQTT 抓包核对 thing/product/{rc_sn}/services 与 services_reply。";
}
return "设备未应答 " + method + "(211001)。请确认网关在线、已订阅 services,并检查 MQTT 连通与 SDKManager 注册。";
}
@Override
public HttpResultResponse requestCloudControlAuth(String sn) {
if (!deviceRedisService.checkDeviceOnline(sn)) {
return HttpResultResponse.error("The gateway is offline.");
}
CloudControlAuthRequest request = new CloudControlAuthRequest()
.setUserId(Optional.ofNullable(SecurityUtils.getUserId()).orElse("cloud"))
.setUserCallsign(Optional.ofNullable(SecurityUtils.getUsername()).orElse("GeoFly"))
.setControlKeys(Collections.singletonList("flight"));
TopicServicesResponse<ServicesReplyData> response =
abstractControlService.cloudControlAuthRequest(SDKManager.getDeviceSDK(sn), request);
ServicesReplyData reply = response.getData();
if (!reply.getResult().isSuccess()) {
return HttpResultResponse.error(reply.getResult());
}
// services_reply 在飞手同意后返回;提示前端/日志
return HttpResultResponse.success().setMessage("云端控制已授权,请继续操作");
}
@Override
public HttpResultResponse releaseCloudControlAuth(String sn) {
if (!deviceRedisService.checkDeviceOnline(sn)) {
return HttpResultResponse.error("The gateway is offline.");
}
CloudControlReleaseRequest request = new CloudControlReleaseRequest()
.setControlKeys(Collections.singletonList("flight"));
TopicServicesResponse<ServicesReplyData> response =
abstractControlService.cloudControlRelease(SDKManager.getDeviceSDK(sn), request);
ServicesReplyData reply = response.getData();
return reply.getResult().isSuccess()
? HttpResultResponse.success()
: HttpResultResponse.error(reply.getResult());
}
private Boolean checkPayloadAuthority(String sn, String payloadIndex) {
Optional<DeviceDTO> dockOpt = deviceRedisService.getDeviceOnline(sn);
if (dockOpt.isEmpty()) {
......
......@@ -23,13 +23,16 @@ import com.dji.sample.wayline.model.param.UpdateJobParam;
import com.dji.sample.wayline.service.IFlightTaskService;
import com.dji.sample.wayline.service.IWaylineJobService;
import com.dji.sample.wayline.service.IWaylineRedisService;
import com.dji.sdk.cloudapi.control.ControlErrorCodeEnum;
import com.dji.sdk.cloudapi.control.ControlMethodEnum;
import com.dji.sdk.cloudapi.control.DrcModeEnterRequest;
import com.dji.sdk.cloudapi.control.DrcModeMqttBroker;
import com.dji.sdk.cloudapi.control.api.AbstractControlService;
import com.dji.sdk.cloudapi.device.DeviceDomainEnum;
import com.dji.sdk.cloudapi.device.DockModeCodeEnum;
import com.dji.sdk.cloudapi.device.OsdDock;
import com.dji.sdk.cloudapi.device.OsdDockDrone;
import com.dji.sdk.cloudapi.device.OsdRcDrone;
import com.dji.sdk.cloudapi.wayline.FlighttaskProgress;
import com.dji.sdk.common.HttpResultResponse;
import com.dji.sdk.common.SDKManager;
......@@ -142,12 +145,36 @@ public class DrcServiceImpl implements IDrcService {
UpdateJobParam.builder().status(WaylineTaskStatusEnum.PAUSE).build());
}
DockModeCodeEnum dockMode = deviceService.getDockMode(dockSn);
Optional<DeviceDTO> dockOpt = deviceRedisService.getDeviceOnline(dockSn);
if (dockOpt.isPresent() && (DockModeCodeEnum.IDLE == dockMode || DockModeCodeEnum.WORKING == dockMode)) {
Optional<OsdDockDrone> deviceOsd = deviceRedisService.getDeviceOsd(dockOpt.get().getChildDeviceSn(), OsdDockDrone.class);
if (dockOpt.isEmpty()) {
throw new RuntimeException("The gateway is offline.");
}
DeviceDTO gateway = dockOpt.get();
// Pilot RC:无 Dock mode / drone_in_dock;校验子机在线后申请云控授权并抢权
if (DeviceDomainEnum.REMOTER_CONTROL == gateway.getDomain()) {
String childSn = gateway.getChildDeviceSn();
if (!StringUtils.hasText(childSn) || !deviceRedisService.checkDeviceOnline(childSn)) {
throw new RuntimeException("Aircraft under RC is offline.");
}
boolean droneAirborne = deviceRedisService.getDeviceOsd(childSn, OsdRcDrone.class)
.map(osd -> osd.getElevation() != null && osd.getElevation() > 0)
.orElse(false);
if (!droneAirborne) {
// 地面也可进 DRC(负载控制),不强制在空;仅打日志
log.info("Pilot RC DRC enter while aircraft may be on ground, sn={}", dockSn);
}
HttpResultResponse auth = controlService.requestCloudControlAuth(dockSn);
if (HttpResultResponse.CODE_SUCCESS != auth.getCode()) {
throw new IllegalArgumentException(auth.getMessage());
}
return;
}
DockModeCodeEnum dockMode = deviceService.getDockMode(dockSn);
if (DockModeCodeEnum.IDLE == dockMode || DockModeCodeEnum.WORKING == dockMode) {
Optional<OsdDockDrone> deviceOsd = deviceRedisService.getDeviceOsd(gateway.getChildDeviceSn(), OsdDockDrone.class);
Optional<OsdDock> dockOsd = deviceRedisService.getDeviceOsd(dockSn, OsdDock.class);
// if (deviceOsd.isEmpty() || deviceOsd.get().getElevation() <= 0) {
if (deviceOsd.isEmpty() || dockOsd.isEmpty() || dockOsd.get().getDroneInDock()) {
throw new RuntimeException("The drone is not in the sky and cannot enter command flight mode.");
}
......@@ -183,7 +210,9 @@ public class DrcServiceImpl implements IDrcService {
TopicServicesResponse<ServicesReplyData> reply = abstractControlService.drcModeEnter(
SDKManager.getDeviceSDK(param.getDockSn()),
new DrcModeEnterRequest()
.setMqttBroker(MqttPropertyConfiguration.getMqttBrokerWithDrc(param.getDockSn() + "-" + System.currentTimeMillis(), param.getDockSn(),
// 设备侧必须用官方 host:port 格式,不能用 Web 的 wss://.../mqtt
.setMqttBroker(MqttPropertyConfiguration.getMqttBrokerWithDrcForDevice(
param.getDockSn() + "-" + System.currentTimeMillis(), param.getDockSn(),
RedisConst.DRC_MODE_ALIVE_SECOND.longValue(),
Map.of(MapKeyConst.ACL, objectMapper.convertValue(JwtAclDTO.builder()
.pub(List.of(subTopic))
......@@ -192,8 +221,15 @@ public class DrcServiceImpl implements IDrcService {
.setHsiFrequency(1).setOsdFrequency(10));
if (!reply.getData().getResult().isSuccess()) {
throw new RuntimeException("SN: " + param.getDockSn() + "; Error:" + reply.getData().getResult() +
"; Failed to enter command flight control mode, please try again later!");
Integer code = reply.getData().getResult().getCode();
String hint = "SN: " + param.getDockSn() + "; Error:" + reply.getData().getResult()
+ "; Failed to enter command flight control mode, please try again later!";
if (code != null && code == ControlErrorCodeEnum.DRC_LINK_REFUSED.getCode()) {
hint += " (514304: 设备连不上 DRC Broker。"
+ "请在 mqtt.DRC 配置 device-host/device-port 为公网可达的 MQTTS,例如 host:8883、device-enable-tls=true;"
+ "官方地址格式为 host:port,不要带 wss://。)";
}
throw new RuntimeException(hint);
}
refreshAcl(param.getDockSn(), param.getClientId(), pubTopic, subTopic);
......@@ -222,8 +258,16 @@ public class DrcServiceImpl implements IDrcService {
TopicServicesResponse<ServicesReplyData> reply =
abstractControlService.drcModeExit(SDKManager.getDeviceSDK(param.getDockSn()));
if (!reply.getData().getResult().isSuccess()) {
throw new RuntimeException("SN: " + param.getDockSn() + "; Error:" +
reply.getData().getResult() + "; Failed to exit command flight control mode, please try again later!");
Integer code = reply.getData().getResult().getCode();
// 官方错误码 514300–514304:网关/网络断开或连接失败,设备侧已无有效 DRC
// 见 docs/dji-cloud-api/ERROR-CODE-NOTES.md 与 cn|en/error-code.md
if (isDrcLinkAlreadyGone(code)) {
log.warn("drc_mode_exit: DRC/gateway link already gone ({}), treat as exited. sn={}",
code, param.getDockSn());
} else {
throw new RuntimeException("SN: " + param.getDockSn() + "; Error:" +
reply.getData().getResult() + "; Failed to exit command flight control mode, please try again later!");
}
}
String jobId = waylineRedisService.getPausedWaylineJobId(param.getDockSn());
......@@ -240,6 +284,31 @@ public class DrcServiceImpl implements IDrcService {
this.delDrcModeInRedis(param.getDockSn());
RedisOpsUtils.del(RedisConst.MQTT_ACL_PREFIX + param.getClientId());
// Pilot:退出 DRC 后释放云端控制授权
Optional<DeviceDTO> gw = deviceRedisService.getDeviceOnline(param.getDockSn());
if (gw.isPresent() && DeviceDomainEnum.REMOTER_CONTROL == gw.get().getDomain()) {
try {
controlService.releaseCloudControlAuth(param.getDockSn());
} catch (Exception e) {
log.warn("release cloud control auth failed for {}: {}", param.getDockSn(), e.getMessage());
}
}
}
/**
* 官方 Error Code 514300–514304:网关异常 / 超时断开 / 证书失败 / 网络断开 / 连接被拒。
* 退出 DRC 时表示链路已不可用,按已退出处理。
*/
private static boolean isDrcLinkAlreadyGone(Integer code) {
if (code == null) {
return false;
}
return code == ControlErrorCodeEnum.DRC_ABNORMAL.getCode()
|| code == ControlErrorCodeEnum.DRC_HEARTBEAT_TIMED_OUT.getCode()
|| code == ControlErrorCodeEnum.DRC_CERTIFICATE_ABNORMAL.getCode()
|| code == ControlErrorCodeEnum.DRC_LINK_LOST.getCode()
|| code == ControlErrorCodeEnum.DRC_LINK_REFUSED.getCode();
}
}
......@@ -5,8 +5,11 @@ import com.dji.sample.control.model.param.DronePayloadParam;
import com.dji.sample.manage.model.dto.DeviceDTO;
import com.dji.sample.manage.service.IDevicePayloadService;
import com.dji.sample.manage.service.IDeviceRedisService;
import com.dji.sdk.cloudapi.device.DeviceDomainEnum;
import com.dji.sdk.cloudapi.device.OsdCamera;
import com.dji.sdk.cloudapi.device.OsdDockDrone;
import com.dji.sdk.cloudapi.device.OsdRcDrone;
import org.springframework.util.StringUtils;
import java.util.Optional;
......@@ -32,16 +35,26 @@ public abstract class PayloadCommandsHandler {
}
public boolean canPublish(String deviceSn) {
Optional<OsdDockDrone> deviceOpt = SpringBeanUtilsTest.getBean(IDeviceRedisService.class)
.getDeviceOsd(deviceSn, OsdDockDrone.class);
if (deviceOpt.isEmpty()) {
throw new RuntimeException("The device is offline.");
IDeviceRedisService redis = SpringBeanUtilsTest.getBean(IDeviceRedisService.class);
Optional<OsdDockDrone> dockDroneOpt = redis.getDeviceOsd(deviceSn, OsdDockDrone.class);
if (dockDroneOpt.isPresent()) {
osdCamera = dockDroneOpt.get().getCameras().stream()
.filter(cam -> param.getPayloadIndex().equals(cam.getPayloadIndex().toString()))
.findAny()
.orElseThrow(() -> new RuntimeException(
"Did not receive osd information about the camera, please check the cache data."));
return true;
}
osdCamera = deviceOpt.get().getCameras().stream()
.filter(osdCamera -> param.getPayloadIndex().equals(osdCamera.getPayloadIndex().toString()))
.findAny()
.orElseThrow(() -> new RuntimeException("Did not receive osd information about the camera, please check the cache data."));
return true;
// Pilot 飞机 OSD 为 OsdRcDrone,通常无 cameras 列表;在线即可下发负载/云台指令
Optional<OsdRcDrone> rcDroneOpt = redis.getDeviceOsd(deviceSn, OsdRcDrone.class);
if (rcDroneOpt.isPresent()) {
osdCamera = null;
return true;
}
throw new RuntimeException("The device is offline.");
}
private String checkDockOnline(String dockSn) {
......@@ -59,7 +72,12 @@ public abstract class PayloadCommandsHandler {
}
}
private void checkAuthority(String deviceSn) {
private void checkAuthority(String dockSn, String deviceSn) {
// Pilot:飞控与负载权不区分,cloud_control_auth(flight) 即可控云台
Optional<DeviceDTO> gatewayOpt = SpringBeanUtilsTest.getBean(IDeviceRedisService.class).getDeviceOnline(dockSn);
if (gatewayOpt.isPresent() && DeviceDomainEnum.REMOTER_CONTROL == gatewayOpt.get().getDomain()) {
return;
}
boolean hasAuthority = SpringBeanUtilsTest.getBean(IDevicePayloadService.class)
.checkAuthorityPayload(deviceSn, param.getPayloadIndex());
if (!hasAuthority) {
......@@ -77,8 +95,11 @@ public abstract class PayloadCommandsHandler {
}
String deviceSn = checkDockOnline(dockSn);
if (!StringUtils.hasText(deviceSn)) {
throw new RuntimeException("The device is offline.");
}
checkDeviceOnline(deviceSn);
checkAuthority(deviceSn);
checkAuthority(dockSn, deviceSn);
if (!canPublish(deviceSn)) {
throw new RuntimeException("The current state of the drone does not support this function, please try again later.");
......
......@@ -53,7 +53,11 @@ public class DeviceRedisServiceImpl implements IDeviceRedisService {
@Override
public <T> Optional<T> getDeviceOsd(String sn, Class<T> clazz) {
return Optional.ofNullable(clazz.cast(RedisOpsUtils.get(RedisConst.OSD_PREFIX + sn)));
Object data = RedisOpsUtils.get(RedisConst.OSD_PREFIX + sn);
if (data == null || !clazz.isInstance(data)) {
return Optional.empty();
}
return Optional.of(clazz.cast(data));
}
@Override
......
......@@ -599,23 +599,37 @@ public class DeviceServiceImpl extends ServiceImpl<IDeviceMapper, DeviceEntity>
}
Page<DeviceEntity> pagination = mapper.selectPage(new Page<>(page, pageSize), wrapper );
List<DeviceDTO> devicesList = pagination.getRecords().stream().map(this::deviceEntityConvertToDTO)
.peek(device -> {
if (StringUtils.hasText(device.getDeviceSn())) {
device.setStatus(deviceRedisService.checkDeviceOnline(device.getDeviceSn()));
}
if (StringUtils.hasText(device.getChildDeviceSn())) {
Optional<DeviceDTO> childOpt = this.getDeviceBySn(device.getChildDeviceSn());
childOpt.ifPresent(child -> {
child.setStatus(deviceRedisService.checkDeviceOnline(child.getDeviceSn()));
child.setWorkspaceName(device.getWorkspaceName());
device.setChildren(child);
});
}
})
.peek(this::enrichDeviceGatewayRelation)
.collect(Collectors.toList());
return new PaginationData<DeviceDTO>(devicesList, new Pagination(pagination.getCurrent(), pagination.getSize(), pagination.getTotal()));
}
/**
* 补充网关/子机关系:机场与 RC 挂 children;飞行器反查 parent_sn(Pilot 模拟机场用 RC SN 下发任务)。
*/
private void enrichDeviceGatewayRelation(DeviceDTO device) {
if (device == null || !StringUtils.hasText(device.getDeviceSn())) {
return;
}
device.setStatus(deviceRedisService.checkDeviceOnline(device.getDeviceSn()));
if (StringUtils.hasText(device.getChildDeviceSn())) {
Optional<DeviceDTO> childOpt = this.getDeviceBySn(device.getChildDeviceSn());
childOpt.ifPresent(child -> {
child.setStatus(deviceRedisService.checkDeviceOnline(child.getDeviceSn()));
child.setWorkspaceName(device.getWorkspaceName());
child.setParentSn(device.getDeviceSn());
device.setChildren(child);
});
}
if (DeviceDomainEnum.DRONE == device.getDomain() && !StringUtils.hasText(device.getParentSn())) {
List<DeviceDTO> parents = this.getDevicesByParams(
DeviceQueryParam.builder().childSn(device.getDeviceSn()).build());
if (!parents.isEmpty()) {
device.setParentSn(parents.get(0).getDeviceSn());
}
}
}
@Override
public PaginationData<DeviceDTO> getDevicesByParam(DeviceSearchParam param,
String workspaceId,
......@@ -692,19 +706,7 @@ public class DeviceServiceImpl extends ServiceImpl<IDeviceMapper, DeviceEntity>
Page<DeviceEntity> pagination = mapper.selectPage(new Page<>(page, pageSize), wrapper );
List<DeviceDTO> devicesList = pagination.getRecords().stream().map(this::deviceEntityConvertToDTO)
.peek(device -> {
if (StringUtils.hasText(device.getDeviceSn())) {
device.setStatus(deviceRedisService.checkDeviceOnline(device.getDeviceSn()));
}
if (StringUtils.hasText(device.getChildDeviceSn())) {
Optional<DeviceDTO> childOpt = this.getDeviceBySn(device.getChildDeviceSn());
childOpt.ifPresent(child -> {
child.setStatus(deviceRedisService.checkDeviceOnline(child.getDeviceSn()));
child.setWorkspaceName(device.getWorkspaceName());
device.setChildren(child);
});
}
})
.peek(this::enrichDeviceGatewayRelation)
.collect(Collectors.toList());
return new PaginationData<>(devicesList, new Pagination(pagination.getCurrent(), pagination.getSize(), pagination.getTotal()));
}
......@@ -783,19 +785,7 @@ public class DeviceServiceImpl extends ServiceImpl<IDeviceMapper, DeviceEntity>
List<DeviceEntity> deviceEntityList = this.list(wrapper);
List<DeviceDTO> devicesList = deviceEntityList.stream().map(this::deviceEntityConvertToDTO)
.peek(device -> {
if (StringUtils.hasText(device.getDeviceSn())) {
device.setStatus(deviceRedisService.checkDeviceOnline(device.getDeviceSn()));
}
if (StringUtils.hasText(device.getChildDeviceSn())) {
Optional<DeviceDTO> childOpt = this.getDeviceBySn(device.getChildDeviceSn());
childOpt.ifPresent(child -> {
child.setStatus(deviceRedisService.checkDeviceOnline(child.getDeviceSn()));
child.setWorkspaceName(device.getWorkspaceName());
device.setChildren(child);
});
}
})
.peek(this::enrichDeviceGatewayRelation)
.collect(Collectors.toList());
return devicesList;
}
......@@ -1037,8 +1027,13 @@ public class DeviceServiceImpl extends ServiceImpl<IDeviceMapper, DeviceEntity>
@Override
public DroneModeCodeEnum getDeviceMode(String deviceSn) {
return deviceRedisService.getDeviceOsd(deviceSn, OsdDockDrone.class)
.map(OsdDockDrone::getModeCode).orElse(DroneModeCodeEnum.DISCONNECTED);
Optional<OsdDockDrone> dockDrone = deviceRedisService.getDeviceOsd(deviceSn, OsdDockDrone.class);
if (dockDrone.isPresent() && dockDrone.get().getModeCode() != null) {
return dockDrone.get().getModeCode();
}
return deviceRedisService.getDeviceOsd(deviceSn, OsdRcDrone.class)
.map(OsdRcDrone::getModeCode)
.orElse(DroneModeCodeEnum.DISCONNECTED);
}
@Override
......@@ -1084,6 +1079,12 @@ public class DeviceServiceImpl extends ServiceImpl<IDeviceMapper, DeviceEntity>
return true;
}
// Pilot RC:无 OsdDock.drc_state,用 Redis 记录是否已进入 DRC
Optional<DeviceDTO> gatewayOpt = deviceRedisService.getDeviceOnline(dockSn);
if (gatewayOpt.isPresent() && DeviceDomainEnum.REMOTER_CONTROL == gatewayOpt.get().getDomain()) {
return RedisOpsUtils.checkExist(RedisConst.DRC_PREFIX + dockSn);
}
return deviceRedisService.getDeviceOsd(dockSn, OsdDock.class)
.map(OsdDock::getDrcState)
.orElse(DrcStateEnum.DISCONNECTED) != DrcStateEnum.DISCONNECTED;
......@@ -1091,11 +1092,14 @@ public class DeviceServiceImpl extends ServiceImpl<IDeviceMapper, DeviceEntity>
@Override
public Boolean checkAuthorityFlight(String gatewaySn) {
return deviceRedisService.getDeviceOnline(gatewaySn).flatMap(gateway ->
Optional.of((DeviceDomainEnum.DOCK == gateway.getDomain()
|| DeviceDomainEnum.REMOTER_CONTROL == gateway.getDomain())
&& ControlSourceEnum.A == gateway.getControlSource()))
.orElse(true);
return deviceRedisService.getDeviceOnline(gatewaySn).map(gateway -> {
// Pilot RC:controlSource 表示遥控器本控 A/B,不能当作云端已获飞控权
if (DeviceDomainEnum.REMOTER_CONTROL == gateway.getDomain()) {
return false;
}
return DeviceDomainEnum.DOCK == gateway.getDomain()
&& ControlSourceEnum.A == gateway.getControlSource();
}).orElse(false);
}
@Override
......
package com.dji.sample.manage.service.impl;
import com.dji.sample.component.websocket.model.BizCodeEnum;
import com.dji.sample.component.websocket.service.IWebSocketMessageService;
import com.dji.sample.manage.model.dto.DeviceDTO;
import com.dji.sample.manage.model.dto.TelemetryDTO;
import com.dji.sample.manage.service.IDeviceRedisService;
import com.dji.sdk.cloudapi.airsense.AirsenseWarning;
import com.dji.sdk.cloudapi.airsense.api.AbstractAirsenseService;
import com.dji.sdk.mqtt.MqttReply;
import com.dji.sdk.mqtt.events.TopicEventsRequest;
import com.dji.sdk.mqtt.events.TopicEventsResponse;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.MessageHeaders;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
/**
* 注册 {@code inboundEventsAirsenseWarning} channel,并处理机场/飞行器 AirSense 告警。
* Demo 未提供实现时,events 路由会因 channel 缺失刷 ERROR。
*/
@Slf4j
@Service
public class SDKAirsenseService extends AbstractAirsenseService {
@Autowired
private IDeviceRedisService deviceRedisService;
@Autowired
private IWebSocketMessageService webSocketMessageService;
@Override
public TopicEventsResponse<MqttReply> airsenseWarning(
TopicEventsRequest<List<AirsenseWarning>> request, MessageHeaders headers) {
String gatewaySn = request.getGateway();
List<AirsenseWarning> warnings = request.getData() == null
? Collections.emptyList() : request.getData();
log.debug("AirSense warning from {}: count={}", gatewaySn, warnings.size());
Optional<DeviceDTO> deviceOpt = deviceRedisService.getDeviceOnline(gatewaySn);
if (deviceOpt.isPresent() && StringUtils.hasText(deviceOpt.get().getWorkspaceId())) {
webSocketMessageService.sendBatch(
deviceOpt.get().getWorkspaceId(),
BizCodeEnum.AIRSENSE_WARNING.getCode(),
TelemetryDTO.<List<AirsenseWarning>>builder()
.sn(gatewaySn)
.host(warnings)
.build());
}
return new TopicEventsResponse<MqttReply>().setData(MqttReply.success());
}
}
......@@ -4,6 +4,8 @@ import cn.hutool.core.date.DateUtil;
import cn.hutool.core.util.NumberUtil;
import cn.hutool.core.util.RandomUtil;
import cn.hutool.json.JSONUtil;
import com.dji.sample.component.redis.RedisConst;
import com.dji.sample.component.redis.RedisOpsUtils;
import com.dji.sample.component.websocket.model.BizCodeEnum;
import com.dji.sample.component.websocket.service.IWebSocketMessageService;
import com.dji.sample.manage.model.dto.DeviceDTO;
......@@ -22,6 +24,7 @@ import com.dji.sdk.config.version.GatewayManager;
import com.dji.sdk.mqtt.MqttReply;
import com.dji.sdk.mqtt.osd.TopicOsdRequest;
import com.dji.sdk.mqtt.state.TopicStateRequest;
import com.dji.sdk.mqtt.state.TopicStateResponse;
import com.dji.sdk.mqtt.status.TopicStatusRequest;
import com.dji.sdk.mqtt.status.TopicStatusResponse;
import lombok.extern.slf4j.Slf4j;
......@@ -387,6 +390,101 @@ public class SDKDeviceService extends AbstractDeviceService {
.build()).collect(Collectors.toList()));
}
/**
* Dock/drone live_status state property. Demo 默认抛异常;此处仅落日志并推送给前端,避免刷 ERROR。
*/
@Override
public void dockLiveStatusUpdate(TopicStateRequest<DockLiveStatus> request, MessageHeaders headers) {
String from = request.getFrom();
DockLiveStatus liveStatus = request.getData();
Optional<DeviceDTO> deviceOpt = deviceRedisService.getDeviceOnline(from);
if (deviceOpt.isEmpty()) {
log.debug("Ignore dock live_status, device offline. SN: {}", from);
return;
}
DeviceDTO device = deviceOpt.get();
if (!StringUtils.hasText(device.getWorkspaceId())) {
log.debug("Ignore dock live_status, device unbound. SN: {}", from);
return;
}
log.debug("Dock live_status from {}: {}", from, liveStatus);
deviceService.pushOsdDataToWeb(device.getWorkspaceId(), BizCodeEnum.DOCK_LIVE_STATUS, from, liveStatus);
}
/**
* RC/drone live_status state property.
*/
@Override
public void rcLiveStatusUpdate(TopicStateRequest<RcLiveStatus> request, MessageHeaders headers) {
String from = request.getFrom();
RcLiveStatus liveStatus = request.getData();
Optional<DeviceDTO> deviceOpt = deviceRedisService.getDeviceOnline(from);
if (deviceOpt.isEmpty()) {
log.debug("Ignore rc live_status, device offline. SN: {}", from);
return;
}
DeviceDTO device = deviceOpt.get();
if (!StringUtils.hasText(device.getWorkspaceId())) {
log.debug("Ignore rc live_status, device unbound. SN: {}", from);
return;
}
log.debug("RC live_status from {}: {}", from, liveStatus);
deviceService.pushOsdDataToWeb(device.getWorkspaceId(), BizCodeEnum.RC_LIVE_STATUS, from, liveStatus);
}
/**
* 机场 / RC / 飞行器上报的 4G Dongle 信息;暂仅落日志,避免 demo 空实现抛异常。
*/
@Override
public TopicStateResponse<MqttReply> dongleInfos(TopicStateRequest<DongleInfos> request, MessageHeaders headers) {
log.debug("dongle_infos from {}: {}", request.getFrom(), request.getData());
return new TopicStateResponse<MqttReply>().setData(MqttReply.success());
}
/**
* 飞行器 WPMZ(航线任务库)版本:缓存到 Redis,并推给前端。
*/
@Override
public void dockWpmzVersionUpdate(TopicStateRequest<DockDroneWpmzVersion> request, MessageHeaders headers) {
String sn = request.getFrom();
String version = request.getData() == null ? null : request.getData().getWpmzVersion();
if (!StringUtils.hasText(version)) {
return;
}
RedisOpsUtils.setWithExpire(RedisConst.WPMZ_VERSION_PREFIX + sn, version, RedisConst.DEVICE_ALIVE_SECOND);
Optional<DeviceDTO> deviceOpt = deviceRedisService.getDeviceOnline(sn);
if (deviceOpt.isEmpty()) {
deviceOpt = deviceService.getDeviceBySn(sn);
}
deviceOpt.filter(d -> StringUtils.hasText(d.getWorkspaceId())).ifPresent(device ->
deviceService.pushOsdDataToWeb(device.getWorkspaceId(), BizCodeEnum.WPMZ_VERSION, sn, request.getData()));
log.info("Updated wpmz_version. SN: {}, version: {}", sn, version);
}
/**
* 飞行器进入当前状态的原因(mode_code_reason):缓存并推送,供前端展示返航/降落等原因。
*/
@Override
public TopicStateResponse<MqttReply> dockDroneModeCodeReason(
TopicStateRequest<DockDroneModeCodeReason> request, MessageHeaders headers) {
String sn = request.getFrom();
DockDroneModeCodeReason data = request.getData();
if (data != null && data.getModeCodeReason() != null) {
RedisOpsUtils.setWithExpire(
RedisConst.MODE_CODE_REASON_PREFIX + sn,
data.getModeCodeReason().getReason(),
RedisConst.DEVICE_ALIVE_SECOND);
Optional<DeviceDTO> deviceOpt = deviceRedisService.getDeviceOnline(sn);
if (deviceOpt.isEmpty()) {
deviceOpt = deviceService.getDeviceBySn(sn);
}
deviceOpt.filter(d -> StringUtils.hasText(d.getWorkspaceId())).ifPresent(device ->
deviceService.pushOsdDataToWeb(device.getWorkspaceId(), BizCodeEnum.MODE_CODE_REASON, sn, data));
log.debug("Updated mode_code_reason. SN: {}, reason: {}", sn, data.getModeCodeReason());
}
return new TopicStateResponse<MqttReply>().setData(MqttReply.success());
}
@Override
public void rcControlSourceUpdate(TopicStateRequest<RcDroneControlSource> request, MessageHeaders headers) {
// If the control source is empty, it will not be processed.
......
......@@ -31,6 +31,8 @@ import com.dji.sdk.cloudapi.wayline.api.AbstractWaylineService;
import com.dji.sdk.common.BaseModel;
import com.dji.sdk.common.HttpResultResponse;
import com.dji.sdk.common.SDKManager;
import com.dji.sdk.exception.CloudSDKErrorEnum;
import com.dji.sdk.exception.CloudSDKException;
import com.dji.sdk.mqtt.MqttReply;
import com.dji.sdk.mqtt.events.TopicEventsRequest;
import com.dji.sdk.mqtt.events.TopicEventsResponse;
......@@ -363,17 +365,80 @@ public class FlightTaskServiceImpl extends AbstractWaylineService implements IFl
}
}
public HttpResultResponse publishOneFlightTask(WaylineJobDTO waylineJob) throws SQLException {
boolean isOnline = deviceRedisService.checkDeviceOnline(waylineJob.getDockSn());
if (!isOnline) {
throw new RuntimeException("Dock is offline.");
/**
* 校验网关可下发航线:机场走 Dock mode_code;Pilot RC「模拟机场」仅要求 RC 与子机在线。
*/
private void assertGatewayReadyForWayline(String gatewaySn) {
Optional<DeviceDTO> gatewayOpt = deviceRedisService.getDeviceOnline(gatewaySn);
if (gatewayOpt.isEmpty()) {
throw new RuntimeException("Gateway is offline.");
}
DeviceDTO gateway = gatewayOpt.get();
if (DeviceDomainEnum.REMOTER_CONTROL == gateway.getDomain()) {
String childSn = gateway.getChildDeviceSn();
if (!StringUtils.hasText(childSn) || !deviceRedisService.checkDeviceOnline(childSn)) {
throw new RuntimeException("Aircraft under RC is offline.");
}
return;
}
DockModeCodeEnum dockModeCodeEnum = deviceRedisService.getDeviceOsd(waylineJob.getDockSn(), OsdDock.class)
DockModeCodeEnum dockModeCodeEnum = deviceRedisService.getDeviceOsd(gatewaySn, OsdDock.class)
.map(OsdDock::getModeCode).orElse(null);
if (dockModeCodeEnum == null || dockModeCodeEnum == DockModeCodeEnum.REMOTE_DEBUGGING) {
throw new RuntimeException("Dock is remote_debugging state, does not support flight task.");
}
}
private boolean isRemoterControlGateway(String gatewaySn) {
return deviceRedisService.getDeviceOnline(gatewaySn)
.map(d -> DeviceDomainEnum.REMOTER_CONTROL == d.getDomain())
.orElse(false);
}
private boolean isMqttNoReply(CloudSDKException e) {
return e.getErrorInfo() != null
&& Objects.equals(CloudSDKErrorEnum.MQTT_PUBLISH_ABNORMAL.getCode(), e.getErrorInfo().getCode());
}
private String mqttNoReplyHint(String gatewaySn, String method) {
if (isRemoterControlGateway(gatewaySn)) {
return "Pilot 遥控器未应答 " + method + "(211001)。"
+ "官方 Pilot 航线管理为 HTTPS 文件库;向 RC 下发的 flighttask_* 需固件支持。"
+ "请确认 dock_sn 为遥控器 SN,并用 MQTT 抓包核对 thing/product/{rc_sn}/services 与 services_reply。";
}
return "设备未应答 " + method + "(211001)。请确认网关在线且已订阅 services,检查 MQTT 连通与 SDKManager 注册。";
}
private void markJobMqttFailed(WaylineJobDTO waylineJob, CloudSDKException e) {
waylineJobService.updateJob(WaylineJobDTO.builder()
.workspaceId(waylineJob.getWorkspaceId())
.jobId(waylineJob.getJobId())
.executeTime(LocalDateTime.now())
.status(WaylineJobStatusEnum.FAILED.getVal())
.completedTime(LocalDateTime.now())
.code(e.getErrorInfo() != null ? e.getErrorInfo().getCode() : CloudSDKErrorEnum.MQTT_PUBLISH_ABNORMAL.getCode())
.build());
}
public HttpResultResponse publishOneFlightTask(WaylineJobDTO waylineJob) throws SQLException {
// 实测:当前 Pilot 固件对 flighttask_* 不回 services_reply(211001)。
// 官方 Pilot 航线为 HTTPS 文件库(App 内 Mission),勿空等 MQTT 超时。
if (isRemoterControlGateway(waylineJob.getDockSn())) {
String msg = "Pilot 上云暂不支持 Web 端 MQTT 航线任务下发(flighttask_prepare)。"
+ "请在 DJI Pilot 2 航线库中执行航线,或改用机场设备创建任务。";
log.warn("Reject wayline job {} for Pilot RC {}: {}", waylineJob.getJobId(), waylineJob.getDockSn(), msg);
waylineJobService.updateJob(WaylineJobDTO.builder()
.workspaceId(waylineJob.getWorkspaceId())
.jobId(waylineJob.getJobId())
.executeTime(LocalDateTime.now())
.status(WaylineJobStatusEnum.FAILED.getVal())
.completedTime(LocalDateTime.now())
.code(CloudSDKErrorEnum.MQTT_PUBLISH_ABNORMAL.getCode())
.build());
return HttpResultResponse.error(msg);
}
assertGatewayReadyForWayline(waylineJob.getDockSn());
boolean isSuccess = this.prepareFlightTask(waylineJob);
if (isSuccess) {
......@@ -462,8 +527,14 @@ public class FlightTaskServiceImpl extends AbstractWaylineService implements IFl
flightTask.setExecutableConditions(waylineJob.getConditions().getExecutableConditions());
}
TopicServicesResponse<ServicesReplyData> serviceReply = abstractWaylineService.flighttaskPrepare(
SDKManager.getDeviceSDK(waylineJob.getDockSn()), flightTask);
TopicServicesResponse<ServicesReplyData> serviceReply;
try {
serviceReply = abstractWaylineService.flighttaskPrepare(
SDKManager.getDeviceSDK(waylineJob.getDockSn()), flightTask);
} catch (CloudSDKException e) {
markJobMqttFailed(waylineJob, e);
throw new RuntimeException(mqttNoReplyHint(waylineJob.getDockSn(), "flighttask_prepare"), e);
}
if (!serviceReply.getData().getResult().isSuccess()) {
log.info("Prepare task ====> Error code: {}", serviceReply.getData().getResult());
waylineJobService.updateJob(WaylineJobDTO.builder()
......@@ -494,8 +565,13 @@ public class FlightTaskServiceImpl extends AbstractWaylineService implements IFl
WaylineJobDTO job = waylineJob.get();
TopicServicesResponse<ServicesReplyData> serviceReply = abstractWaylineService.flighttaskExecute(
SDKManager.getDeviceSDK(job.getDockSn()), new FlighttaskExecuteRequest().setFlightId(jobId));
TopicServicesResponse<ServicesReplyData> serviceReply;
try {
serviceReply = abstractWaylineService.flighttaskExecute(
SDKManager.getDeviceSDK(job.getDockSn()), new FlighttaskExecuteRequest().setFlightId(jobId));
} catch (CloudSDKException e) {
throw new RuntimeException(mqttNoReplyHint(job.getDockSn(), "flighttask_execute"), e);
}
if (!serviceReply.getData().getResult().isSuccess()) {
log.info("Execute job ====> Error: {}", serviceReply.getData().getResult());
waylineJobService.updateJob(WaylineJobDTO.builder()
......@@ -546,11 +622,24 @@ public class FlightTaskServiceImpl extends AbstractWaylineService implements IFl
throw new RuntimeException("Dock is offline.");
}
TopicServicesResponse<ServicesReplyData> serviceReply = abstractWaylineService.flighttaskUndo(SDKManager.getDeviceSDK(dockSn),
new FlighttaskUndoRequest().setFlightIds(jobIds));
if (!serviceReply.getData().getResult().isSuccess()) {
log.info("Cancel job ====> Error: {}", serviceReply.getData().getResult());
throw new RuntimeException("Failed to cancel the wayline job of " + dockSn);
boolean mqttOk = false;
try {
TopicServicesResponse<ServicesReplyData> serviceReply = abstractWaylineService.flighttaskUndo(
SDKManager.getDeviceSDK(dockSn),
new FlighttaskUndoRequest().setFlightIds(jobIds));
if (!serviceReply.getData().getResult().isSuccess()) {
log.info("Cancel job ====> Error: {}", serviceReply.getData().getResult());
throw new RuntimeException("Failed to cancel the wayline job of " + dockSn);
}
mqttOk = true;
} catch (CloudSDKException e) {
// Pilot RC:PENDING 任务常因设备从不应答 flighttask_*;允许仅落库取消,避免删任务卡死
if (isMqttNoReply(e) && isRemoterControlGateway(dockSn)) {
log.warn("flighttask_undo no reply from Pilot RC {}, soft-cancel jobs {} in DB only. {}",
dockSn, jobIds, e.getMessage());
} else {
throw new RuntimeException(mqttNoReplyHint(dockSn, "flighttask_undo"), e);
}
}
for (String jobId : jobIds) {
......@@ -562,6 +651,9 @@ public class FlightTaskServiceImpl extends AbstractWaylineService implements IFl
.build());
RedisOpsUtils.zRemove(RedisConst.WAYLINE_JOB_TIMED_EXECUTE, workspaceId + RedisConst.DELIMITER + dockSn + RedisConst.DELIMITER + jobId);
}
if (!mqttOk) {
log.info("Soft-cancelled {} job(s) for Pilot RC {} without device undo ack", jobIds.size(), dockSn);
}
}
......@@ -772,15 +864,7 @@ public class FlightTaskServiceImpl extends AbstractWaylineService implements IFl
public HttpResultResponse publishOneInFlightFlightTask(WaylineJobDTO waylineJob) throws SQLException {
boolean isOnline = deviceRedisService.checkDeviceOnline(waylineJob.getDockSn());
if (!isOnline) {
throw new RuntimeException("Dock is offline.");
}
DockModeCodeEnum dockModeCodeEnum = deviceRedisService.getDeviceOsd(waylineJob.getDockSn(), OsdDock.class)
.map(OsdDock::getModeCode).orElse(null);
if (dockModeCodeEnum == null || dockModeCodeEnum == DockModeCodeEnum.REMOTE_DEBUGGING) {
throw new RuntimeException("Dock is remote_debugging state, does not support flight task.");
}
assertGatewayReadyForWayline(waylineJob.getDockSn());
boolean isSuccess = this.prepareInFlightTask(waylineJob);
if (!isSuccess) {
......
......@@ -72,7 +72,17 @@ public class SDKWaylineService extends AbstractWaylineService {
@Override
public TopicEventsResponse<MqttReply> deviceExitHomingNotify(TopicEventsRequest<DeviceExitHomingNotify> request, MessageHeaders headers) {
return super.deviceExitHomingNotify(request, headers);
log.debug("device_exit_homing_notify from {}: {}", request.getGateway(), request.getData());
return new TopicEventsResponse<MqttReply>().setData(MqttReply.success());
}
@Override
public TopicEventsResponse<MqttReply> returnHomeInfo(TopicEventsRequest<ReturnHomeInfo> request, MessageHeaders headers) {
log.debug("return_home_info from {}: flightId={}, lastPointType={}",
request.getGateway(),
request.getData() == null ? null : request.getData().getFlightId(),
request.getData() == null ? null : request.getData().getLastPointType());
return new TopicEventsResponse<MqttReply>().setData(MqttReply.success());
}
@Override
......
......@@ -95,16 +95,13 @@ mqtt:
protocol: WSS
host: geofly.geotwin.cn
port: 443
# protocol: WSS # @see com.dji.sample.component.mqtt.model.MqttProtocolEnum
# host: emqx-broker
# host: 192.168.32.90
# port: 8083
# host: geofly-dev.geotwin.cc
# host: geofly.geotwin.cc
# port: 443
path: /mqtt
username: JavaServer
password: 123456
# 下发给机场/Pilot:54418 明文 MQTT → enable-tls=false;MQTTS 8883 则改为 true
device-host: geofly.geotwin.cn
device-port: 54418
device-enable-tls: false
cloud-sdk:
mqtt:
......
server:
port: 6789
my:
file:
base-path: /app/filePath
spring:
main:
allow-bean-definition-overriding: true
application:
# name: cloud-api-sample
name: geoair-api
datasource:
druid:
type: com.alibaba.druid.pool.DruidDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
url: jdbc:mysql://10.0.8.80:3306/cloud_sample?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Shanghai
# url: jdbc:mysql://localhost:3306/cloud_sample?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Shanghai
username: root
password: root
initial-size: 10
min-idle: 10
max-active: 20
max-wait: 60000
redis:
# host: geoair_redis
# host: cloud_api_sample_redis
# 深圳
host: 10.0.8.80
port: 6379
database: 0
username: # if you enable
password:
lettuce:
pool:
max-active: 8
max-idle: 8
min-idle: 0
servlet:
multipart:
max-file-size: 2GB
max-request-size: 2GB
# 邮箱
mail:
host: smtpdm.aliyun.com #阿里云发送服务器地址
port: 25 #端口号
username: geofly@stmp-geofly.geotwin.cc #发送人地址
# password: ENC(Grg2n2TYzgJv9zpwufsf37ndTe+M1cYk) #密码
password: GeoSys2326 #密码
jwt:
issuer: DJI
subject: CloudApiSample
secret: CloudApiSample
age: 86400
mqtt:
# @see com.dji.sample.component.mqtt.model.MqttUseEnum
# BASIC parameters are required.
NET:
# protocol: MQTT # @see com.dji.sample.component.mqtt.model.MqttProtocolEnum
# host: 183.11.236.162
# port: 54418
# username: JavaServer
# password: 123456
# client-id: 123456
protocol: WSS
host: geofly.geotwin.cn
port: 443
path: /mqtt
username: JavaServer
password: 123456
BASIC:
protocol: MQTT # @see com.dji.sample.component.mqtt.model.MqttProtocolEnum
# 深圳
# host: 183.11.236.162
# port: 54418
# host: 203.186.109.106
# host: emqx-broker
# host: 192.168.32.90
# port: 44418
# host: 203.186.109.106
# port: 54941
host: 10.0.8.80
port: 1883
username: JavaServer
password: 123456
client-id: 123456
# If the protocol is ws/wss, this value is required.
path:
DRC:
# 深圳
protocol: WSS
host: geofly.geotwin.cn
port: 443
path: /mqtt
username: JavaServer
password: 123456
# 下发给机场/Pilot:54418 明文 MQTT → enable-tls=false;MQTTS 8883 则改为 true
device-host: geofly.geotwin.cn
device-port: 54418
device-enable-tls: false
cloud-sdk:
mqtt:
# Topics that need to be subscribed when initially connecting to mqtt, multiple topics are divided by ",".
inbound-topic: sys/product/+/status,thing/product/+/requests,thing/product/+/osd,/ai_info
url:
manage:
prefix: manage
version: /api/v1
map:
prefix: map
version: /api/v1
media:
prefix: media
version: /api/v1
wayline:
prefix: wayline
version: /api/v1
storage:
prefix: storage
version: /api/v1
control:
prefix: control
version: /api/v1
psdk:
prefix: psdk
version: /api/v1
# Tutorial: https://www.alibabacloud.com/help/en/object-storage-service/latest/use-a-temporary-credential-provided-by-sts-to-access-oss
#oss:
# enable: false
# provider: ALIYUN # @see com.dji.sample.component.OssConfiguration.model.enums.OssTypeEnum
# endpoint: https://oss-cn-hangzhou.aliyuncs.com
# access-key: Please enter your access key.
# secret-key: Please enter your secret key.
# expire: 3600
# region: Please enter your oss region. # cn-hangzhou
# role-session-name: cloudApi
# role-arn: Please enter your role arn. # acs:ram::123456789:role/stsrole
# bucket: Please enter your bucket name.
# object-dir-prefix: Please enter a folder name.
#oss:
# enable: true
# provider: aws
# endpoint: https://s3.us-east-1.amazonaws.com
# access-key:
# secret-key:
# expire: 3600
# region: us-east-1
# role-session-name: cloudApi
# role-arn:
# bucket: cloudapi-bucket
# object-dir-prefix: wayline
# VersityGW:provider=versitygw(OssTypeEnum.VERSITYGW)。下发给 Pilot/机场时会映射为 minio。
# sts-mode=temp:Admin API 创建短期用户(需开启 VersityGW 多租户 IAM);static:直接下发长期 AK/SK。
# 若仍用 MinIO,改为 provider: minio 即可。
oss:
enable: true
provider: versitygw
# 香港
# endpoint: https://gt7-oss.geotwin.cc
# 深圳 / VersityGW
endpoint: https://gt-oss-dev.geotwin.cn
access-key: minioadmin
secret-key: minioadmin
bucket: gtfly
expire: 86400 # 24 * 3600 24小时
region: us-east-1 # us-east-1
object-dir-prefix: wayline
sts-mode: temp
sts-user-role: admin
logging:
level:
com.dji: info
file:
name: logs/cloud-api-sample.log
ntp:
server:
host: Google.mzr.me
# To create a license for an application: https://developer.dji.com/user/apps/#all
cloud-api:
app:
id: 163614
key: 09e2fc9ac2ade8607d07826b3c1a822
license: fvwK+dVYFBUwxGUxXS3EglmsuTbyFA6400LGuqLrmebHmnhLL9U5DPTI0l4Wu49Xc6fdmNjP6e7ZzHqEhhwloVvhSn4tKMc9wAIsTLlHz1r/j7YbHo5xDX2LHLJOss1e8ud8Gl2lC3gfo9Dsc+qZiCfBX/uIeeZUgk/rR8m1erE=
# 旧配置
# id: 158488
# key: 15e3972c29ef16b9eda5e3415d943d7
# license: el2GS5tHPIEjXazwGZC1LkIEQ3B43U7uQ9x30ateJAACX1miOSOdjKLn0clBJTLKO3Fd/EnrlNnVwoPOmYvn9K6skalfLgY3bBv5I6sJq0nzaIGpdCHelNA/81GCvb9j1sWdN3adRIk/v5GtmSswauqElmgBGN/r37yp7pi5b5E=
# id: 158487
# key: 11de4b6e406e51bf78c75c4cabefe09
# license: MQzV5RtcqtuZuPlP6MAq0/1xF7gD2ujuN+0tmh+rDzOb7xECeIcRoyGItLcicTGWu3OkunX4XB8/7RcUp5yjyZhyvjHuqhkRDxSwfToidln8mRJ9UvltcSgKRCYCYnR/Oa6u3u0bs2azz20t6zCKmfdAIY0GzHI7YqgmOHowBXU=
livestream:
url:
# It is recommended to use a program to create Token. https://github.com/AgoraIO/Tools/blob/master/DynamicKey/AgoraDynamicKey/java/src/main/java/io/agora/media/RtcTokenBuilder2.java
agora:
channel: Please enter the agora channel.
token: Please enter the agora temporary token.
uid: 654321
# RTMP Note: This IP is the address of the streaming server. If you want to see livestream on web page, you need to convert the RTMP stream to WebRTC stream.
rtmp:
# url: rtmp://203.186.109.106:44424/live/ # Example: 'rtmp://192.168.1.1/live/'
# 深圳
# url: rtmp://183.11.236.162:54424/live/ # Example: 'rtmp://192.168.1.1/live/'
url: rtmp://183.11.236.162:54460/live/ # Example: 'rtmp://192.168.1.1/live/'
rtsp:
username: Please enter the username.
password: Please enter the password.
port: 8554
# GB28181 Note:If you don't know what these parameters mean, you can go to Pilot2 and select the GB28181 page in the cloud platform. Where the parameters same as these parameters.
gb28181:
serverIP: Please enter the server ip.
serverPort: 0
serverID: Please enter the server id.
agentID: Please enter the agent id.
agentPassword: Please enter the agent password.
localPort: 0
channel: Please enter the channel.
# Webrtc: Only supports using whip standard
whip:
url: Please enter the rtmp access address. # Example:http://192.168.1.1:1985/rtc/v1/whip/?app=live&stream=
uom:
# 数据源 Integer
# 1:无人机系统直报的数据
# 2:无人机制造商自建的无人机运行服务系统代报的数据
# 3:无人机机体上加装的单独数据模块代报的数据
# 4:采集设备采集代报的广播数据
source: 2
# 识别信息报送平台
# 报送平台的名称,如“DJI FLY”
platform: "GEO FLY"
# 报送程序版本
# 程序版本号,以第一次完成报送对接记为“version1.0”,每次升级报送程序累加版本号后缀
programVersion: "version1.0"
# 上报平台 API 地址 https://uom.receive.caacic.cn/addFlightRoute
# url: https://218.189.32.212:8080/addFlightRoute
url: https://220.232.168.7:8080/addFlightRoute
appID: GEOSYS
appKey: geosys_uas
......@@ -98,16 +98,13 @@ mqtt:
protocol: WSS
host: geofly.geotwin.cn
port: 443
# protocol: WSS # @see com.dji.sample.component.mqtt.model.MqttProtocolEnum
# host: emqx-broker
# host: 192.168.32.90
# port: 8083
# host: geofly-dev.geotwin.cc
# host: geofly.geotwin.cc
# port: 443
path: /mqtt
username: JavaServer
password: 123456
# 下发给机场/Pilot:54418 明文 MQTT → enable-tls=false;MQTTS 8883 则改为 true
device-host: geofly.geotwin.cn
device-port: 54418
device-enable-tls: false
cloud-sdk:
mqtt:
......@@ -185,7 +182,7 @@ oss:
logging:
level:
com.dji: debug
com.dji: info
file:
name: logs/cloud-api-sample.log
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment