Browse Source

MQTT backward compatibility adaptor: init commit

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

2
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java

@ -54,6 +54,8 @@ public class ProtoTransportPayloadConfiguration implements TransportPayloadTypeC
private String deviceRpcRequestProtoSchema;
private String deviceRpcResponseProtoSchema;
private boolean enableCompatibilityWithOtherPayloadFormats;
@Override
public TransportPayloadType getTransportPayloadType() {
return TransportPayloadType.PROTOBUF;

16
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -967,6 +967,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
log.trace("[{}] Received get attributes response", sessionId);
String topicBase;
MqttTransportAdaptor adaptor;
boolean useBackupAdaptorByDefault = false;
switch (attrReqTopicType) {
case V2:
adaptor = deviceSessionCtx.getPayloadAdaptor();
@ -983,10 +984,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
default:
adaptor = deviceSessionCtx.getPayloadAdaptor();
topicBase = MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_TOPIC_PREFIX;
useBackupAdaptorByDefault = true;
break;
}
try {
adaptor.convertToPublish(deviceSessionCtx, response, topicBase).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
adaptor.convertToPublish(deviceSessionCtx, response, topicBase, useBackupAdaptorByDefault).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
} catch (Exception e) {
log.trace("[{}] Failed to convert device attributes response to MQTT msg", sessionId, e);
}
@ -998,6 +1000,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
log.info("[{}] : attrSubTopicType: {}", notification.toString(), attrSubTopicType);
String topic;
MqttTransportAdaptor adaptor;
boolean useBackupAdaptorByDefault = false;
switch (attrSubTopicType) {
case V2:
adaptor = deviceSessionCtx.getPayloadAdaptor();
@ -1014,10 +1017,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
default:
adaptor = deviceSessionCtx.getPayloadAdaptor();
topic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC;
useBackupAdaptorByDefault = true;
break;
}
try {
adaptor.convertToPublish(deviceSessionCtx, notification, topic).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
adaptor.convertToPublish(deviceSessionCtx, notification, topic, useBackupAdaptorByDefault).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
} catch (Exception e) {
log.trace("[{}] Failed to convert device attributes update to MQTT msg", sessionId, e);
}
@ -1034,6 +1038,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
log.info("[{}] Received RPC command to device", sessionId);
String baseTopic;
MqttTransportAdaptor adaptor;
boolean useBackupAdaptorByDefault = false;
switch (rpcSubTopicType) {
case V2:
adaptor = deviceSessionCtx.getPayloadAdaptor();
@ -1050,10 +1055,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
default:
adaptor = deviceSessionCtx.getPayloadAdaptor();
baseTopic = MqttTopics.DEVICE_RPC_REQUESTS_TOPIC;
useBackupAdaptorByDefault = true;
break;
}
try {
adaptor.convertToPublish(deviceSessionCtx, rpcRequest, baseTopic).ifPresent(payload -> {
adaptor.convertToPublish(deviceSessionCtx, rpcRequest, baseTopic, useBackupAdaptorByDefault).ifPresent(payload -> {
int msgId = ((MqttPublishMessage) payload).variableHeader().packetId();
if (isAckExpected(payload)) {
rpcAwaitingAck.put(msgId, rpcRequest);
@ -1090,6 +1096,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
log.trace("[{}] Received RPC response from server", sessionId);
String baseTopic;
MqttTransportAdaptor adaptor;
boolean useBackupAdaptorByDefault = false;
switch (toServerRpcSubTopicType) {
case V2:
adaptor = deviceSessionCtx.getPayloadAdaptor();
@ -1106,10 +1113,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
default:
adaptor = deviceSessionCtx.getPayloadAdaptor();
baseTopic = MqttTopics.DEVICE_RPC_RESPONSE_TOPIC;
useBackupAdaptorByDefault = true;
break;
}
try {
adaptor.convertToPublish(deviceSessionCtx, rpcResponse, baseTopic).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
adaptor.convertToPublish(deviceSessionCtx, rpcResponse, baseTopic, useBackupAdaptorByDefault).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
} catch (Exception e) {
log.trace("[{}] Failed to convert device RPC command to MQTT msg", sessionId, e);
}

139
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java

@ -0,0 +1,139 @@
/**
* 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.adaptors;
import io.netty.handler.codec.mqtt.MqttMessage;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.mqtt.session.MqttDeviceAwareSessionContext;
import java.util.Optional;
@Data
@AllArgsConstructor
@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;
@Override
public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
try {
return main.convertToPostTelemetry(ctx, inbound);
} catch (AdaptorException e) {
log.trace(BACKWARD_COMPATIBILITY_ENABLED + "post telemetry request msg using {} ...", backup.getClass().getSimpleName());
return backup.convertToPostTelemetry(ctx, inbound);
}
}
@Override
public TransportProtos.PostAttributeMsg convertToPostAttributes(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
try {
return main.convertToPostAttributes(ctx, inbound);
} catch (AdaptorException e) {
log.trace(BACKWARD_COMPATIBILITY_ENABLED + "post attributes request msg using {} ...", backup.getClass().getSimpleName());
return backup.convertToPostAttributes(ctx, inbound);
}
}
@Override
public TransportProtos.GetAttributeRequestMsg convertToGetAttributes(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound, String topicBase) throws AdaptorException {
try {
return main.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);
}
}
@Override
public TransportProtos.ToDeviceRpcResponseMsg convertToDeviceRpcResponse(MqttDeviceAwareSessionContext ctx, MqttPublishMessage mqttMsg, String topicBase) throws AdaptorException {
try {
return main.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);
}
}
@Override
public TransportProtos.ToServerRpcRequestMsg convertToServerRpcRequest(MqttDeviceAwareSessionContext ctx, MqttPublishMessage mqttMsg, String topicBase) throws AdaptorException {
try {
return main.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);
}
}
@Override
public TransportProtos.ClaimDeviceMsg convertToClaimDevice(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
try {
return main.convertToClaimDevice(ctx, inbound);
} catch (AdaptorException e) {
log.trace(BACKWARD_COMPATIBILITY_ENABLED + "claim device request msg using {} ...", backup.getClass().getSimpleName());
return backup.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);
}
@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);
}
@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);
}
@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);
}
}

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) throws AdaptorException {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) 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) {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) {
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) {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) {
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) {
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) {
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,23 +60,29 @@ public interface MqttTransportAdaptor {
ClaimDeviceMsg convertToClaimDevice(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, GetAttributeResponseMsg responseMsg, String topicBase) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, GetAttributeResponseMsg responseMsg, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException;
Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, GetAttributeResponseMsg responseMsg) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, AttributeUpdateNotificationMsg notificationMsg, String topic) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, AttributeUpdateNotificationMsg notificationMsg, String topic, boolean useBackupAdaptorByDefault) throws AdaptorException;
Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToDeviceRpcRequestMsg rpcRequest, String topicBase) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToDeviceRpcRequestMsg rpcRequest, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException;
Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase) throws AdaptorException;
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase, boolean useBackupAdaptorByDefault) throws AdaptorException;
ProvisionDeviceRequestMsg convertToProvisionRequestMsg(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException;
default ProvisionDeviceRequestMsg convertToProvisionRequestMsg(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
// method is never used in BackwardCompatibilityAdaptor
return null;
}
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ProvisionDeviceResponseMsg provisionResponse) throws AdaptorException;
default Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, ProvisionDeviceResponseMsg provisionResponse) throws AdaptorException {
// method is never used in BackwardCompatibilityAdaptor
return Optional.empty();
}
Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, byte[] firmwareChunk, String requestId, int chunk, OtaPackageType firmwareType) throws AdaptorException;

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

13
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.adaptors.BackwardCompatibilityAdaptor;
import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor;
import org.thingsboard.server.transport.mqtt.util.MqttTopicFilter;
import org.thingsboard.server.transport.mqtt.util.MqttTopicFilterFactory;
@ -76,6 +77,7 @@ 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;
@ -103,7 +105,15 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
}
public MqttTransportAdaptor getPayloadAdaptor() {
return payloadType.equals(TransportPayloadType.JSON) ? context.getJsonMqttAdaptor() : context.getProtoMqttAdaptor();
if (payloadType.equals(TransportPayloadType.JSON)) {
return context.getJsonMqttAdaptor();
} else {
if (payloadFormatsCompatipilityEnabled) {
return new BackwardCompatibilityAdaptor(context.getProtoMqttAdaptor(), context.getJsonMqttAdaptor());
} else {
return context.getProtoMqttAdaptor();
}
}
}
public boolean isJsonPayloadType() {
@ -162,6 +172,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
attributesDynamicMessageDescriptor = protoTransportPayloadConfig.getAttributesDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceAttributesProtoSchema());
rpcResponseDynamicMessageDescriptor = protoTransportPayloadConfig.getRpcResponseDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceRpcResponseProtoSchema());
rpcRequestDynamicMessageBuilder = protoTransportPayloadConfig.getRpcRequestDynamicMessageBuilder(protoTransportPayloadConfig.getDeviceRpcRequestProtoSchema());
payloadFormatsCompatipilityEnabled = protoTransportPayloadConfig.isEnableCompatibilityWithOtherPayloadFormats();
}
public void addToQueue(MqttMessage msg) {

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

@ -73,6 +73,12 @@
{{ 'device-profile.mqtt-payload-type-required' | translate }}
</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>
<div class="tb-hint" innerHTML="{{ 'device-profile.mqtt-enable-compatibility-with-other-payload-formats-hint' | translate }}"></div>
</div>
<div *ngIf="protoPayloadType" fxLayout="column">
<mat-form-field fxFlex>
<mat-label translate>device-profile.telemetry-proto-schema</mat-label>

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

@ -97,7 +97,8 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control
deviceTelemetryProtoSchema: [defaultTelemetrySchema, Validators.required],
deviceAttributesProtoSchema: [defaultAttributesSchema, Validators.required],
deviceRpcRequestProtoSchema: [defaultRpcRequestSchema, Validators.required],
deviceRpcResponseProtoSchema: [defaultRpcResponseSchema, Validators.required]
deviceRpcResponseProtoSchema: [defaultRpcResponseSchema, Validators.required],
enableCompatibilityWithOtherPayloadFormats: [false, Validators.required]
})
}, {validator: this.uniqueDeviceTopicValidator}
);
@ -156,7 +157,8 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control
deviceTelemetryProtoSchema: defaultTelemetrySchema,
deviceAttributesProtoSchema: defaultAttributesSchema,
deviceRpcRequestProtoSchema: defaultRpcRequestSchema,
deviceRpcResponseProtoSchema: defaultRpcResponseSchema
deviceRpcResponseProtoSchema: defaultRpcResponseSchema,
enableCompatibilityWithOtherPayloadFormats: false,
}, {emitEvent: false});
}
if (type === TransportPayloadType.PROTOBUF && !this.disabled) {
@ -164,11 +166,13 @@ 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});
} 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});
}
}

6
ui-ngx/src/app/shared/models/device.models.ts

@ -247,6 +247,7 @@ export interface MqttDeviceProfileTransportConfiguration {
deviceAttributesTopic?: string;
transportPayloadTypeConfiguration?: {
transportPayloadType?: TransportPayloadType;
enableCompatibilityWithOtherPayloadFormats?: boolean;
};
[key: string]: any;
}
@ -359,7 +360,10 @@ export function createDeviceProfileTransportConfiguration(type: DeviceTransportT
const mqttTransportConfiguration: MqttDeviceProfileTransportConfiguration = {
deviceTelemetryTopic: 'v1/devices/me/telemetry',
deviceAttributesTopic: 'v1/devices/me/attributes',
transportPayloadTypeConfiguration: {transportPayloadType: TransportPayloadType.JSON}
transportPayloadTypeConfiguration: {
transportPayloadType: TransportPayloadType.JSON,
enableCompatibilityWithOtherPayloadFormats: false
}
};
transportConfiguration = {...mqttTransportConfiguration, type: DeviceTransportType.MQTT};
break;

2
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -1096,6 +1096,8 @@
"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.",
"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