Browse Source

update logic after review

pull/5441/head
ShvaykaD 5 years ago
parent
commit
75c1185d18
  1. 3
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java
  2. 44
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  3. 20
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TopicType.java
  4. 70
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java
  5. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java
  6. 18
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java
  7. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
  8. 60
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
  9. 12
      ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html
  10. 54
      ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts
  11. 6
      ui-ngx/src/app/shared/models/device.models.ts
  12. 6
      ui-ngx/src/assets/locale/locale.constant-en_US.json

3
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() {

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

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

70
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<MqttMessage> 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<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException {
return main.convertToGatewayPublish(ctx, deviceName, responseMsg);
}
@Override
public Optional<MqttMessage> 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<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException {
return main.convertToGatewayPublish(ctx, deviceName, notificationMsg);
}
@Override
public Optional<MqttMessage> 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<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException {
return main.convertToGatewayPublish(ctx, deviceName, rpcRequest);
}
@Override
public Optional<MqttMessage> 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<MqttMessage> 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);
}
}

8
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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException {
public Optional<MqttMessage> 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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) {
public Optional<MqttMessage> 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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) {
public Optional<MqttMessage> 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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase) {
return Optional.of(createMqttPublishMsg(ctx, topicBase + rpcResponse.getRequestId(), JsonConverter.toJson(rpcResponse)));
}

18
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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException;
default Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, GetAttributeResponseMsg responseMsg, String topicBase) throws AdaptorException {
return Optional.empty();
}
Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, GetAttributeResponseMsg responseMsg) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) throws AdaptorException;
default Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, AttributeUpdateNotificationMsg notificationMsg, String topic) throws AdaptorException {
return Optional.empty();
}
Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException;
default Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToDeviceRpcRequestMsg rpcRequest, String topicBase) throws AdaptorException {
return Optional.empty();
}
Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException;
default Optional<MqttMessage> 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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ProvisionDeviceResponseMsg provisionResponse) throws AdaptorException {
// method is never used in BackwardCompatibilityAdaptor
return Optional.empty();
}

8
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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException {
public Optional<MqttMessage> 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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException {
public Optional<MqttMessage> 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<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase) {
return Optional.of(createMqttPublishMsg(ctx, topicBase + rpcResponse.getRequestId(), rpcResponse.toByteArray()));
}
@Override
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) {
return Optional.of(createMqttPublishMsg(ctx, topic, notificationMsg.toByteArray()));
}

60
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<MqttTopicMatcher, Integer> 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) {

12
ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html

@ -74,10 +74,16 @@
</mat-error>
</mat-form-field>
<div *ngIf="protoPayloadType" style="padding-bottom: 20px">
<mat-checkbox formControlName="enableCompatibilityWithOtherPayloadFormats">
{{ 'device-profile.mqtt-enable-compatibility-with-other-payload-formats' | translate }}
<mat-checkbox formControlName="enableCompatibilityWithJsonPayloadFormat">
{{ 'device-profile.mqtt-enable-compatibility-with-json-payload-format' | translate }}
</mat-checkbox>
<div class="tb-hint" innerHTML="{{ 'device-profile.mqtt-enable-compatibility-with-other-payload-formats-hint' | translate }}"></div>
<div class="tb-hint" innerHTML="{{ 'device-profile.mqtt-enable-compatibility-with-json-payload-format-hint' | translate }}"></div>
<div *ngIf="compatibilityWithJsonPayloadFormatEnabled">
<mat-checkbox formControlName="useJsonPayloadFormatForDefaultDownlinkTopics">
{{ 'device-profile.mqtt-use-json-format-for-default-downlink-topics' | translate }}
</mat-checkbox>
<div class="tb-hint" innerHTML="{{ 'device-profile.mqtt-use-json-format-for-default-downlink-topics-hint' | translate }}"></div>
</div>
</div>
<div *ngIf="protoPayloadType" fxLayout="column">
<mat-form-field fxFlex>

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

6
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};

6
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: <b>v1/devices/me/attributes/response/$request_id</b>, <b>v1/devices/me/attributes</b>, <b>v1/devices/me/rpc/request/$request_id</b>, <b>v1/devices/me/rpc/response/$request_id</b>. This setting does not impact attribute and rpc subscriptions sent using new (v2) topics: <b>v2/a/res/$request_id</b>, <b>v2/a</b>, <b>v2/r/req/$request_id</b>, <b>v2/r/res/$request_id</b>. <b>$request_id</b> 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",

Loading…
Cancel
Save