diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java index 1e45ff3ba6..08a0cdd610 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java @@ -54,7 +54,8 @@ public class ProtoTransportPayloadConfiguration implements TransportPayloadTypeC private String deviceRpcRequestProtoSchema; private String deviceRpcResponseProtoSchema; - private boolean enableCompatibilityWithOtherPayloadFormats; + private boolean enableCompatibilityWithJsonPayloadFormat; + private boolean useJsonPayloadFormatForDefaultDownlinkTopics; @Override public TransportPayloadType getTransportPayloadType() { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 3d04ecfdc9..f9ddbf1058 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -966,29 +966,23 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg response) { log.trace("[{}] Received get attributes response", sessionId); String topicBase; - MqttTransportAdaptor adaptor; - boolean useBackupAdaptorByDefault = false; + MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(attrReqTopicType); switch (attrReqTopicType) { case V2: - adaptor = deviceSessionCtx.getPayloadAdaptor(); topicBase = MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_SHORT_TOPIC_PREFIX; break; case V2_JSON: - adaptor = context.getJsonMqttAdaptor(); topicBase = MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_SHORT_JSON_TOPIC_PREFIX; break; case V2_PROTO: - adaptor = context.getProtoMqttAdaptor(); topicBase = MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_SHORT_PROTO_TOPIC_PREFIX; break; default: - adaptor = deviceSessionCtx.getPayloadAdaptor(); topicBase = MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_TOPIC_PREFIX; - useBackupAdaptorByDefault = true; break; } try { - adaptor.convertToPublish(deviceSessionCtx, response, topicBase, useBackupAdaptorByDefault).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); + adaptor.convertToPublish(deviceSessionCtx, response, topicBase).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); } catch (Exception e) { log.trace("[{}] Failed to convert device attributes response to MQTT msg", sessionId, e); } @@ -999,29 +993,23 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement log.trace("[{}] Received attributes update notification to device", sessionId); log.info("[{}] : attrSubTopicType: {}", notification.toString(), attrSubTopicType); String topic; - MqttTransportAdaptor adaptor; - boolean useBackupAdaptorByDefault = false; + MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(attrSubTopicType); switch (attrSubTopicType) { case V2: - adaptor = deviceSessionCtx.getPayloadAdaptor(); topic = MqttTopics.DEVICE_ATTRIBUTES_SHORT_TOPIC; break; case V2_JSON: - adaptor = context.getJsonMqttAdaptor(); topic = MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC; break; case V2_PROTO: - adaptor = context.getProtoMqttAdaptor(); topic = MqttTopics.DEVICE_ATTRIBUTES_SHORT_PROTO_TOPIC; break; default: - adaptor = deviceSessionCtx.getPayloadAdaptor(); topic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; - useBackupAdaptorByDefault = true; break; } try { - adaptor.convertToPublish(deviceSessionCtx, notification, topic, useBackupAdaptorByDefault).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); + adaptor.convertToPublish(deviceSessionCtx, notification, topic).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); } catch (Exception e) { log.trace("[{}] Failed to convert device attributes update to MQTT msg", sessionId, e); } @@ -1037,29 +1025,23 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { log.info("[{}] Received RPC command to device", sessionId); String baseTopic; - MqttTransportAdaptor adaptor; - boolean useBackupAdaptorByDefault = false; + MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(rpcSubTopicType); switch (rpcSubTopicType) { case V2: - adaptor = deviceSessionCtx.getPayloadAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_REQUESTS_SHORT_TOPIC; break; case V2_JSON: - adaptor = context.getJsonMqttAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_REQUESTS_SHORT_JSON_TOPIC; break; case V2_PROTO: - adaptor = context.getProtoMqttAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_REQUESTS_SHORT_PROTO_TOPIC; break; default: - adaptor = deviceSessionCtx.getPayloadAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_REQUESTS_TOPIC; - useBackupAdaptorByDefault = true; break; } try { - adaptor.convertToPublish(deviceSessionCtx, rpcRequest, baseTopic, useBackupAdaptorByDefault).ifPresent(payload -> { + adaptor.convertToPublish(deviceSessionCtx, rpcRequest, baseTopic).ifPresent(payload -> { int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); if (isAckExpected(payload)) { rpcAwaitingAck.put(msgId, rpcRequest); @@ -1095,29 +1077,23 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg rpcResponse) { log.trace("[{}] Received RPC response from server", sessionId); String baseTopic; - MqttTransportAdaptor adaptor; - boolean useBackupAdaptorByDefault = false; + MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(toServerRpcSubTopicType); switch (toServerRpcSubTopicType) { case V2: - adaptor = deviceSessionCtx.getPayloadAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_RESPONSE_SHORT_TOPIC; break; case V2_JSON: - adaptor = context.getJsonMqttAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_RESPONSE_SHORT_JSON_TOPIC; break; case V2_PROTO: - adaptor = context.getProtoMqttAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_RESPONSE_SHORT_PROTO_TOPIC; break; default: - adaptor = deviceSessionCtx.getPayloadAdaptor(); baseTopic = MqttTopics.DEVICE_RPC_RESPONSE_TOPIC; - useBackupAdaptorByDefault = true; break; } try { - adaptor.convertToPublish(deviceSessionCtx, rpcResponse, baseTopic, useBackupAdaptorByDefault).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); + adaptor.convertToPublish(deviceSessionCtx, rpcResponse, baseTopic).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); } catch (Exception e) { log.trace("[{}] Failed to convert device RPC command to MQTT msg", sessionId, e); } @@ -1141,8 +1117,4 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement deviceSessionCtx.onDeviceUpdate(sessionInfo, device, deviceProfileOpt); } - private enum TopicType { - V1, V2, V2_JSON, V2_PROTO - } - } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TopicType.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TopicType.java new file mode 100644 index 0000000000..e2f427261e --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TopicType.java @@ -0,0 +1,20 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt; + +public enum TopicType { + V1, V2, V2_JSON, V2_PROTO +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java index 4b3a3e25b5..141332f7a7 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java @@ -32,108 +32,86 @@ import java.util.Optional; @Slf4j public class BackwardCompatibilityAdaptor implements MqttTransportAdaptor { - private static final String BACKWARD_COMPATIBILITY_ENABLED = "Other payload formats compatibility enabled! Trying to convert "; - - private MqttTransportAdaptor main; - private MqttTransportAdaptor backup; + private MqttTransportAdaptor protoAdaptor; + private MqttTransportAdaptor jsonAdaptor; @Override public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { try { - return main.convertToPostTelemetry(ctx, inbound); + return protoAdaptor.convertToPostTelemetry(ctx, inbound); } catch (AdaptorException e) { - log.trace(BACKWARD_COMPATIBILITY_ENABLED + "post telemetry request msg using {} ...", backup.getClass().getSimpleName()); - return backup.convertToPostTelemetry(ctx, inbound); + log.trace("failed to process post telemetry request msg {} due to: ", inbound, e); + return jsonAdaptor.convertToPostTelemetry(ctx, inbound); } } @Override public TransportProtos.PostAttributeMsg convertToPostAttributes(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { try { - return main.convertToPostAttributes(ctx, inbound); + return protoAdaptor.convertToPostAttributes(ctx, inbound); } catch (AdaptorException e) { - log.trace(BACKWARD_COMPATIBILITY_ENABLED + "post attributes request msg using {} ...", backup.getClass().getSimpleName()); - return backup.convertToPostAttributes(ctx, inbound); + log.trace("failed to process post attributes request msg {} due to: ", inbound, e); + return jsonAdaptor.convertToPostAttributes(ctx, inbound); } } @Override public TransportProtos.GetAttributeRequestMsg convertToGetAttributes(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound, String topicBase) throws AdaptorException { try { - return main.convertToGetAttributes(ctx, inbound, topicBase); + return protoAdaptor.convertToGetAttributes(ctx, inbound, topicBase); } catch (AdaptorException e) { - log.trace(BACKWARD_COMPATIBILITY_ENABLED + "get attributes request msg using {} ...", backup.getClass().getSimpleName()); - return backup.convertToGetAttributes(ctx, inbound, topicBase); + log.trace("failed to process get attributes request msg {} due to: ", inbound, e); + return jsonAdaptor.convertToGetAttributes(ctx, inbound, topicBase); } } @Override public TransportProtos.ToDeviceRpcResponseMsg convertToDeviceRpcResponse(MqttDeviceAwareSessionContext ctx, MqttPublishMessage mqttMsg, String topicBase) throws AdaptorException { try { - return main.convertToDeviceRpcResponse(ctx, mqttMsg, topicBase); + return protoAdaptor.convertToDeviceRpcResponse(ctx, mqttMsg, topicBase); } catch (AdaptorException e) { - log.trace(BACKWARD_COMPATIBILITY_ENABLED + "to device rpc response msg using {} ...", backup.getClass().getSimpleName()); - return backup.convertToDeviceRpcResponse(ctx, mqttMsg, topicBase); + log.trace("failed to process to device rpc response msg {} due to: ", mqttMsg, e); + return jsonAdaptor.convertToDeviceRpcResponse(ctx, mqttMsg, topicBase); } } @Override public TransportProtos.ToServerRpcRequestMsg convertToServerRpcRequest(MqttDeviceAwareSessionContext ctx, MqttPublishMessage mqttMsg, String topicBase) throws AdaptorException { try { - return main.convertToServerRpcRequest(ctx, mqttMsg, topicBase); + return protoAdaptor.convertToServerRpcRequest(ctx, mqttMsg, topicBase); } catch (AdaptorException e) { - log.trace(BACKWARD_COMPATIBILITY_ENABLED + "to server rpc request msg using {} ...", backup.getClass().getSimpleName()); - return backup.convertToServerRpcRequest(ctx, mqttMsg, topicBase); + log.trace("failed to process to server rpc request msg {} due to: ", mqttMsg, e); + return jsonAdaptor.convertToServerRpcRequest(ctx, mqttMsg, topicBase); } } @Override public TransportProtos.ClaimDeviceMsg convertToClaimDevice(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { try { - return main.convertToClaimDevice(ctx, inbound); + return protoAdaptor.convertToClaimDevice(ctx, inbound); } catch (AdaptorException e) { - log.trace(BACKWARD_COMPATIBILITY_ENABLED + "claim device request msg using {} ...", backup.getClass().getSimpleName()); - return backup.convertToClaimDevice(ctx, inbound); + log.trace("failed to process claim device request msg {} due to: ", inbound, e); + return jsonAdaptor.convertToClaimDevice(ctx, inbound); } } - @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException { - return useBackupAdaptorByDefault ? backup.convertToPublish(ctx, responseMsg, topicBase, false) : main.convertToPublish(ctx, responseMsg, topicBase, false); - } - @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { - return main.convertToGatewayPublish(ctx, deviceName, responseMsg); - } - - @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) throws AdaptorException { - return useBackupAdaptorByDefault ? backup.convertToPublish(ctx, notificationMsg, topic, false) : main.convertToPublish(ctx, notificationMsg, topic, false); + return protoAdaptor.convertToGatewayPublish(ctx, deviceName, responseMsg); } @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException { - return main.convertToGatewayPublish(ctx, deviceName, notificationMsg); - } - - @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException { - return useBackupAdaptorByDefault ? backup.convertToPublish(ctx, rpcRequest, topicBase, false) : main.convertToPublish(ctx, rpcRequest, topicBase, false); + return protoAdaptor.convertToGatewayPublish(ctx, deviceName, notificationMsg); } @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException { - return main.convertToGatewayPublish(ctx, deviceName, rpcRequest); - } - - @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException { - return useBackupAdaptorByDefault ? backup.convertToPublish(ctx, rpcResponse, topicBase, false) : main.convertToPublish(ctx, rpcResponse, topicBase, false); + return protoAdaptor.convertToGatewayPublish(ctx, deviceName, rpcRequest); } @Override public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, byte[] firmwareChunk, String requestId, int chunk, OtaPackageType firmwareType) throws AdaptorException { - return main.convertToPublish(ctx, firmwareChunk, requestId, chunk, firmwareType); + return protoAdaptor.convertToPublish(ctx, firmwareChunk, requestId, chunk, firmwareType); } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index 68d80538c1..dc2b74a8ea 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -117,7 +117,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase) throws AdaptorException { return processConvertFromAttributeResponseMsg(ctx, responseMsg, topicBase); } @@ -127,7 +127,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) { return Optional.of(createMqttPublishMsg(ctx, topic, JsonConverter.toJson(notificationMsg))); } @@ -138,7 +138,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase) { return Optional.of(createMqttPublishMsg(ctx, topicBase + rpcRequest.getRequestId(), JsonConverter.toJson(rpcRequest, false))); } @@ -148,7 +148,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase) { return Optional.of(createMqttPublishMsg(ctx, topicBase + rpcResponse.getRequestId(), JsonConverter.toJson(rpcResponse))); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java index f030dd9fca..1543b5c6c7 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java @@ -60,27 +60,33 @@ public interface MqttTransportAdaptor { ClaimDeviceMsg convertToClaimDevice(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException; - Optional convertToPublish(MqttDeviceAwareSessionContext ctx, GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException; + default Optional convertToPublish(MqttDeviceAwareSessionContext ctx, GetAttributeResponseMsg responseMsg, String topicBase) throws AdaptorException { + return Optional.empty(); + } Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, GetAttributeResponseMsg responseMsg) throws AdaptorException; - Optional convertToPublish(MqttDeviceAwareSessionContext ctx, AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) throws AdaptorException; + default Optional convertToPublish(MqttDeviceAwareSessionContext ctx, AttributeUpdateNotificationMsg notificationMsg, String topic) throws AdaptorException { + return Optional.empty(); + } Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException; - Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException; + default Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ToDeviceRpcRequestMsg rpcRequest, String topicBase) throws AdaptorException { + return Optional.empty(); + } Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException; - Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException; + default Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase) throws AdaptorException { + return Optional.empty(); + } default ProvisionDeviceRequestMsg convertToProvisionRequestMsg(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { - // method is never used in BackwardCompatibilityAdaptor return null; } default Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ProvisionDeviceResponseMsg provisionResponse) throws AdaptorException { - // method is never used in BackwardCompatibilityAdaptor return Optional.empty(); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java index 7f9b326f9d..4fa47b367d 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java @@ -135,7 +135,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase) throws AdaptorException { if (!StringUtils.isEmpty(responseMsg.getError())) { throw new AdaptorException(responseMsg.getError()); } else { @@ -148,7 +148,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase) throws AdaptorException { DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx; DynamicMessage.Builder rpcRequestDynamicMessageBuilder = deviceSessionCtx.getRpcRequestDynamicMessageBuilder(); if (rpcRequestDynamicMessageBuilder == null) { @@ -159,12 +159,12 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase) { return Optional.of(createMqttPublishMsg(ctx, topicBase + rpcResponse.getRequestId(), rpcResponse.toByteArray())); } @Override - public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) { + public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) { return Optional.of(createMqttPublishMsg(ctx, topic, notificationMsg.toByteArray())); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index cc65da1008..585587d6f0 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadCo import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.mqtt.MqttTransportContext; +import org.thingsboard.server.transport.mqtt.TopicType; import org.thingsboard.server.transport.mqtt.adaptors.BackwardCompatibilityAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; import org.thingsboard.server.transport.mqtt.util.MqttTopicFilter; @@ -39,7 +40,6 @@ import org.thingsboard.server.transport.mqtt.util.MqttTopicFilterFactory; import java.util.Collection; import java.util.Collections; -import java.util.Queue; import java.util.UUID; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentMap; @@ -77,11 +77,14 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { private volatile MqttTopicFilter telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); private volatile MqttTopicFilter attributesTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile TransportPayloadType payloadType = TransportPayloadType.JSON; - private volatile boolean payloadFormatsCompatipilityEnabled; private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor; private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor; private volatile Descriptors.Descriptor rpcResponseDynamicMessageDescriptor; private volatile DynamicMessage.Builder rpcRequestDynamicMessageBuilder; + private volatile MqttTransportAdaptor adaptor; + private volatile boolean jsonPayloadFormatCompatibilityEnabled; + private volatile boolean useJsonPayloadFormatForDefaultDownlinkTopics; + @Getter @Setter @@ -90,6 +93,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { public DeviceSessionCtx(UUID sessionId, ConcurrentMap mqttQoSMap, MqttTransportContext context) { super(sessionId, mqttQoSMap); this.context = context; + this.adaptor = context.getJsonMqttAdaptor(); } public int nextMsgId() { @@ -105,15 +109,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { } public MqttTransportAdaptor getPayloadAdaptor() { - if (payloadType.equals(TransportPayloadType.JSON)) { - return context.getJsonMqttAdaptor(); - } else { - if (payloadFormatsCompatipilityEnabled) { - return new BackwardCompatibilityAdaptor(context.getProtoMqttAdaptor(), context.getJsonMqttAdaptor()); - } else { - return context.getProtoMqttAdaptor(); - } - } + return adaptor; } public boolean isJsonPayloadType() { @@ -140,12 +136,14 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { public void setDeviceProfile(DeviceProfile deviceProfile) { super.setDeviceProfile(deviceProfile); updateTopicFilters(deviceProfile); + updateAdaptor(); } @Override public void onDeviceProfileUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceProfile deviceProfile) { super.onDeviceProfileUpdate(sessionInfo, deviceProfile); updateTopicFilters(deviceProfile); + updateAdaptor(); } private void updateTopicFilters(DeviceProfile deviceProfile) { @@ -158,7 +156,10 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); attributesTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); if (TransportPayloadType.PROTOBUF.equals(payloadType)) { - updateDynamicMessageDescriptors(transportPayloadTypeConfiguration); + ProtoTransportPayloadConfiguration protoTransportPayloadConfig = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; + updateDynamicMessageDescriptors(protoTransportPayloadConfig); + jsonPayloadFormatCompatibilityEnabled = protoTransportPayloadConfig.isEnableCompatibilityWithJsonPayloadFormat(); + useJsonPayloadFormatForDefaultDownlinkTopics = protoTransportPayloadConfig.isUseJsonPayloadFormatForDefaultDownlinkTopics(); } } else { telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); @@ -166,13 +167,42 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { } } - private void updateDynamicMessageDescriptors(TransportPayloadTypeConfiguration transportPayloadTypeConfiguration) { - ProtoTransportPayloadConfiguration protoTransportPayloadConfig = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; + private void updateDynamicMessageDescriptors(ProtoTransportPayloadConfiguration protoTransportPayloadConfig) { telemetryDynamicMessageDescriptor = protoTransportPayloadConfig.getTelemetryDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceTelemetryProtoSchema()); attributesDynamicMessageDescriptor = protoTransportPayloadConfig.getAttributesDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceAttributesProtoSchema()); rpcResponseDynamicMessageDescriptor = protoTransportPayloadConfig.getRpcResponseDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceRpcResponseProtoSchema()); rpcRequestDynamicMessageBuilder = protoTransportPayloadConfig.getRpcRequestDynamicMessageBuilder(protoTransportPayloadConfig.getDeviceRpcRequestProtoSchema()); - payloadFormatsCompatipilityEnabled = protoTransportPayloadConfig.isEnableCompatibilityWithOtherPayloadFormats(); + } + + public MqttTransportAdaptor getAdaptor(TopicType topicType) { + switch (topicType) { + case V2: + return getProfileAdaptor(); + case V2_JSON: + return context.getJsonMqttAdaptor(); + case V2_PROTO: + return context.getProtoMqttAdaptor(); + default: + return useJsonPayloadFormatForDefaultDownlinkTopics ? context.getJsonMqttAdaptor() : getProfileAdaptor(); + } + } + + private MqttTransportAdaptor getProfileAdaptor() { + return isJsonPayloadType() ? context.getJsonMqttAdaptor() : context.getProtoMqttAdaptor(); + } + + private void updateAdaptor() { + if (isJsonPayloadType()) { + adaptor = context.getJsonMqttAdaptor(); + jsonPayloadFormatCompatibilityEnabled = false; + useJsonPayloadFormatForDefaultDownlinkTopics = false; + } else { + if (jsonPayloadFormatCompatibilityEnabled) { + adaptor = new BackwardCompatibilityAdaptor(context.getProtoMqttAdaptor(), context.getJsonMqttAdaptor()); + } else { + adaptor = context.getProtoMqttAdaptor(); + } + } } public void addToQueue(MqttMessage msg) { diff --git a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html index a72718bcb5..5f49053c42 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html @@ -74,10 +74,16 @@
- - {{ 'device-profile.mqtt-enable-compatibility-with-other-payload-formats' | translate }} + + {{ 'device-profile.mqtt-enable-compatibility-with-json-payload-format' | translate }} -
+
+
+ + {{ 'device-profile.mqtt-use-json-format-for-default-downlink-topics' | translate }} + +
+
diff --git a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts index 4e97737880..a8da7cc4e7 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts +++ b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts @@ -98,7 +98,8 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control deviceAttributesProtoSchema: [defaultAttributesSchema, Validators.required], deviceRpcRequestProtoSchema: [defaultRpcRequestSchema, Validators.required], deviceRpcResponseProtoSchema: [defaultRpcResponseSchema, Validators.required], - enableCompatibilityWithOtherPayloadFormats: [false, Validators.required] + enableCompatibilityWithJsonPayloadFormat: [false, Validators.required], + useJsonPayloadFormatForDefaultDownlinkTopics: [false, Validators.required] }) }, {validator: this.uniqueDeviceTopicValidator} ); @@ -107,6 +108,14 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control ).subscribe(payloadType => { this.updateTransportPayloadBasedControls(payloadType, true); }); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('transportPayloadTypeConfiguration.enableCompatibilityWithJsonPayloadFormat') + .valueChanges.pipe(takeUntil(this.destroy$) + ).subscribe(compatibilityWithJsonPayloadFormatEnabled => { + if (!compatibilityWithJsonPayloadFormatEnabled) { + this.mqttDeviceProfileTransportConfigurationFormGroup.get('transportPayloadTypeConfiguration.useJsonPayloadFormatForDefaultDownlinkTopics') + .patchValue(false, {emitEvent: false}); + } + }); this.mqttDeviceProfileTransportConfigurationFormGroup.valueChanges.pipe( takeUntil(this.destroy$) ).subscribe(() => { @@ -133,6 +142,10 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control return transportPayloadType === TransportPayloadType.PROTOBUF; } + get compatibilityWithJsonPayloadFormatEnabled(): boolean { + return this.mqttDeviceProfileTransportConfigurationFormGroup.get('transportPayloadTypeConfiguration.enableCompatibilityWithJsonPayloadFormat').value; + } + writeValue(value: MqttDeviceProfileTransportConfiguration | null): void { if (isDefinedAndNotNull(value)) { this.mqttDeviceProfileTransportConfigurationFormGroup.patchValue(value, {emitEvent: false}); @@ -149,6 +162,36 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control this.propagateChange(configuration); } + private updateTest(type: TransportPayloadType, forceUpdated = false) { + const transportPayloadTypeForm = this.mqttDeviceProfileTransportConfigurationFormGroup + .get('transportPayloadTypeConfiguration') as FormGroup; + if (forceUpdated) { + transportPayloadTypeForm.patchValue({ + deviceTelemetryProtoSchema: defaultTelemetrySchema, + deviceAttributesProtoSchema: defaultAttributesSchema, + deviceRpcRequestProtoSchema: defaultRpcRequestSchema, + deviceRpcResponseProtoSchema: defaultRpcResponseSchema, + enableCompatibilityWithJsonPayloadFormat: false, + useJsonPayloadFormatForDefaultDownlinkTopics: false + }, {emitEvent: false}); + } + if (type === TransportPayloadType.PROTOBUF && !this.disabled) { + transportPayloadTypeForm.get('deviceTelemetryProtoSchema').enable({emitEvent: false}); + transportPayloadTypeForm.get('deviceAttributesProtoSchema').enable({emitEvent: false}); + transportPayloadTypeForm.get('deviceRpcRequestProtoSchema').enable({emitEvent: false}); + transportPayloadTypeForm.get('deviceRpcResponseProtoSchema').enable({emitEvent: false}); + transportPayloadTypeForm.get('enableCompatibilityWithJsonPayloadFormat').enable({emitEvent: false}); + transportPayloadTypeForm.get('useJsonPayloadFormatForDefaultDownlinkTopics').enable({emitEvent: false}); + } else { + transportPayloadTypeForm.get('deviceTelemetryProtoSchema').disable({emitEvent: false}); + transportPayloadTypeForm.get('deviceAttributesProtoSchema').disable({emitEvent: false}); + transportPayloadTypeForm.get('deviceRpcRequestProtoSchema').disable({emitEvent: false}); + transportPayloadTypeForm.get('deviceRpcResponseProtoSchema').disable({emitEvent: false}); + transportPayloadTypeForm.get('enableCompatibilityWithJsonPayloadFormat').disable({emitEvent: false}); + transportPayloadTypeForm.get('useJsonPayloadFormatForDefaultDownlinkTopics').disable({emitEvent: false}); + } + } + private updateTransportPayloadBasedControls(type: TransportPayloadType, forceUpdated = false) { const transportPayloadTypeForm = this.mqttDeviceProfileTransportConfigurationFormGroup .get('transportPayloadTypeConfiguration') as FormGroup; @@ -158,7 +201,8 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control deviceAttributesProtoSchema: defaultAttributesSchema, deviceRpcRequestProtoSchema: defaultRpcRequestSchema, deviceRpcResponseProtoSchema: defaultRpcResponseSchema, - enableCompatibilityWithOtherPayloadFormats: false, + enableCompatibilityWithJsonPayloadFormat: false, + useJsonPayloadFormatForDefaultDownlinkTopics: false }, {emitEvent: false}); } if (type === TransportPayloadType.PROTOBUF && !this.disabled) { @@ -166,13 +210,15 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control transportPayloadTypeForm.get('deviceAttributesProtoSchema').enable({emitEvent: false}); transportPayloadTypeForm.get('deviceRpcRequestProtoSchema').enable({emitEvent: false}); transportPayloadTypeForm.get('deviceRpcResponseProtoSchema').enable({emitEvent: false}); - transportPayloadTypeForm.get('enableCompatibilityWithOtherPayloadFormats').enable({emitEvent: false}); + transportPayloadTypeForm.get('enableCompatibilityWithJsonPayloadFormat').enable({emitEvent: false}); + transportPayloadTypeForm.get('useJsonPayloadFormatForDefaultDownlinkTopics').enable({emitEvent: false}); } else { transportPayloadTypeForm.get('deviceTelemetryProtoSchema').disable({emitEvent: false}); transportPayloadTypeForm.get('deviceAttributesProtoSchema').disable({emitEvent: false}); transportPayloadTypeForm.get('deviceRpcRequestProtoSchema').disable({emitEvent: false}); transportPayloadTypeForm.get('deviceRpcResponseProtoSchema').disable({emitEvent: false}); - transportPayloadTypeForm.get('enableCompatibilityWithOtherPayloadFormats').disable({emitEvent: false}); + transportPayloadTypeForm.get('enableCompatibilityWithJsonPayloadFormat').disable({emitEvent: false}); + transportPayloadTypeForm.get('useJsonPayloadFormatForDefaultDownlinkTopics').disable({emitEvent: false}); } } diff --git a/ui-ngx/src/app/shared/models/device.models.ts b/ui-ngx/src/app/shared/models/device.models.ts index 806e47b6ce..7ebf7bde47 100644 --- a/ui-ngx/src/app/shared/models/device.models.ts +++ b/ui-ngx/src/app/shared/models/device.models.ts @@ -247,7 +247,8 @@ export interface MqttDeviceProfileTransportConfiguration { deviceAttributesTopic?: string; transportPayloadTypeConfiguration?: { transportPayloadType?: TransportPayloadType; - enableCompatibilityWithOtherPayloadFormats?: boolean; + enableCompatibilityWithJsonPayloadFormat?: boolean; + useJsonPayloadFormatForDefaultDownlinkTopics?: boolean; }; [key: string]: any; } @@ -362,7 +363,8 @@ export function createDeviceProfileTransportConfiguration(type: DeviceTransportT deviceAttributesTopic: 'v1/devices/me/attributes', transportPayloadTypeConfiguration: { transportPayloadType: TransportPayloadType.JSON, - enableCompatibilityWithOtherPayloadFormats: false + enableCompatibilityWithJsonPayloadFormat: false, + useJsonPayloadFormatForDefaultDownlinkTopics: false, } }; transportConfiguration = {...mqttTransportConfiguration, type: DeviceTransportType.MQTT}; diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 40e5976ed0..85a685aa4b 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1096,8 +1096,10 @@ "mqtt-device-payload-type": "MQTT device payload", "mqtt-device-payload-type-json": "JSON", "mqtt-device-payload-type-proto": "Protobuf", - "mqtt-enable-compatibility-with-other-payload-formats": "Enable compatibility with other payload formats.", - "mqtt-enable-compatibility-with-other-payload-formats-hint": "When enabled, the platform will use a specified payload format by default. If parsing fails, the platform will attempt to use other available formats. Useful for backward compatibility during firmware updates. For example, the initial release of the firmware uses JSON, while the new release uses Protobuf. During the process of firmware update for the fleet of devices, it is required to support both Protobuf and JSON simultaneously. The compatibility mode introduces slight performance degradation, so it is recommended to disable this mode once all devices are updated.", + "mqtt-enable-compatibility-with-json-payload-format": "Enable compatibility with other payload formats.", + "mqtt-enable-compatibility-with-json-payload-format-hint": "When enabled, the platform will use a Protobuf payload format by default. If parsing fails, the platform will attempt to use JSON payload format. Useful for backward compatibility during firmware updates. For example, the initial release of the firmware uses Json, while the new release uses Protobuf. During the process of firmware update for the fleet of devices, it is required to support both Protobuf and JSON simultaneously. The compatibility mode introduces slight performance degradation, so it is recommended to disable this mode once all devices are updated.", + "mqtt-use-json-format-for-default-downlink-topics": "Use Json format for default downlink topics", + "mqtt-use-json-format-for-default-downlink-topics-hint": "When enabled, the platform will use Json payload format to push attributes and RPC via the following topics: v1/devices/me/attributes/response/$request_id, v1/devices/me/attributes, v1/devices/me/rpc/request/$request_id, v1/devices/me/rpc/response/$request_id. This setting does not impact attribute and rpc subscriptions sent using new (v2) topics: v2/a/res/$request_id, v2/a, v2/r/req/$request_id, v2/r/res/$request_id. $request_id is an integer request identifier.", "snmp-add-mapping": "Add SNMP mapping", "snmp-mapping-not-configured": "No mapping for OID to timeseries/telemetry configured", "snmp-timseries-or-attribute-name": "Timeseries/attribute name for mapping",