diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java index 48c5ab6978..fef67ccc72 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java @@ -99,12 +99,10 @@ public class DeviceProfileController extends BaseController { DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); - if (transportConfiguration instanceof MqttDeviceProfileTransportConfiguration) { - if (transportConfiguration instanceof MqttProtoDeviceProfileTransportConfiguration) { - MqttProtoDeviceProfileTransportConfiguration protoTransportConfiguration = (MqttProtoDeviceProfileTransportConfiguration) transportConfiguration; - if (protoTransportConfiguration.getTransportPayloadType().equals(TransportPayloadType.PROTOBUF)) + if (transportConfiguration instanceof MqttProtoDeviceProfileTransportConfiguration) { + MqttProtoDeviceProfileTransportConfiguration protoTransportConfiguration = (MqttProtoDeviceProfileTransportConfiguration) transportConfiguration; + if (protoTransportConfiguration.getTransportPayloadType().equals(TransportPayloadType.PROTOBUF)) checkProtoSchemas(protoTransportConfiguration); - } } DeviceProfile savedDeviceProfile = checkNotNull(deviceProfileService.saveDeviceProfile(deviceProfile)); diff --git a/common/data/pom.xml b/common/data/pom.xml index 00e53ff2bb..6d32a3df53 100644 --- a/common/data/pom.xml +++ b/common/data/pom.xml @@ -71,10 +71,6 @@ java-driver-core test - - org.springframework - spring-web - com.squareup.wire wire-schema diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java index 3395e398e8..720e81b8fd 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java @@ -38,11 +38,11 @@ import org.thingsboard.server.common.data.DeviceTransportType; @JsonDeserialize(using = MqttTransportConfigurationDeserializer.class) public abstract class MqttDeviceProfileTransportConfiguration implements DeviceProfileTransportConfiguration { + public abstract TransportPayloadType getTransportPayloadType(); + protected String deviceTelemetryTopic = MqttTopics.DEVICE_TELEMETRY_TOPIC; protected String deviceAttributesTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; - public abstract TransportPayloadType getTransportPayloadType(); - @Override public DeviceTransportType getType() { return DeviceTransportType.MQTT; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttProtoDeviceProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttProtoDeviceProfileTransportConfiguration.java index c3c9fe263a..06e8d57c8e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttProtoDeviceProfileTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttProtoDeviceProfileTransportConfiguration.java @@ -22,6 +22,7 @@ import com.github.os72.protobuf.dynamic.EnumDefinition; import com.github.os72.protobuf.dynamic.MessageDefinition; import com.google.protobuf.Descriptors; import com.google.protobuf.DynamicMessage; +import com.squareup.wire.schema.Field; import com.squareup.wire.schema.Location; import com.squareup.wire.schema.internal.parser.EnumConstantElement; import com.squareup.wire.schema.internal.parser.EnumElement; @@ -33,8 +34,6 @@ import com.squareup.wire.schema.internal.parser.TypeElement; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.extern.slf4j.Slf4j; -import org.springframework.util.CollectionUtils; -import org.springframework.util.StringUtils; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.TransportPayloadType; @@ -73,7 +72,7 @@ public class MqttProtoDeviceProfileTransportConfiguration extends MqttDeviceProf throw new IllegalArgumentException("Failed to parse: " + schemaName + " due to: " + e.getMessage()); } List types = protoFileElement.getTypes(); - if (!CollectionUtils.isEmpty(types)) { + if (!types.isEmpty()) { if (types.stream().noneMatch(typeElement -> typeElement instanceof MessageElement)) { throw new IllegalArgumentException("Invalid " + schemaName + " provided! At least one Message definition should exists!"); } @@ -86,7 +85,7 @@ public class MqttProtoDeviceProfileTransportConfiguration extends MqttDeviceProf ProtoFileElement protoFileElement = getTransportProtoSchema(protoSchema); DynamicSchema.Builder schemaBuilder = DynamicSchema.newBuilder(); schemaBuilder.setName(schemaName); - schemaBuilder.setPackage(!StringUtils.isEmpty(protoFileElement.getPackageName()) ? + schemaBuilder.setPackage(!isEmptyStr(protoFileElement.getPackageName()) ? protoFileElement.getPackageName() : schemaName.toLowerCase()); List types = protoFileElement.getTypes(); @@ -101,11 +100,11 @@ public class MqttProtoDeviceProfileTransportConfiguration extends MqttDeviceProf .map(typeElement -> (MessageElement) typeElement) .collect(Collectors.toList()); - if (!CollectionUtils.isEmpty(enumTypes)) { + if (!enumTypes.isEmpty()) { enumTypes.forEach(enumElement -> { List enumElementTypeConstants = enumElement.getConstants(); EnumDefinition.Builder enumDefinitionBuilder = EnumDefinition.newBuilder(enumElement.getName()); - if (!CollectionUtils.isEmpty(enumElementTypeConstants)) { + if (!enumElementTypeConstants.isEmpty()) { enumElementTypeConstants.forEach(constantElement -> enumDefinitionBuilder.addValue(constantElement.getName(), constantElement.getTag())); } EnumDefinition enumDefinition = enumDefinitionBuilder.build(); @@ -113,12 +112,21 @@ public class MqttProtoDeviceProfileTransportConfiguration extends MqttDeviceProf }); } - if (!CollectionUtils.isEmpty(messageTypes)) { + if (!messageTypes.isEmpty()) { messageTypes.forEach(messageElement -> { List messageElementFields = messageElement.getFields(); MessageDefinition.Builder messageDefinitionBuilder = MessageDefinition.newBuilder(messageElement.getName()); - if (!CollectionUtils.isEmpty(messageElementFields)) { - messageElementFields.forEach(fieldElement -> messageDefinitionBuilder.addField(fieldElement.getType(), fieldElement.getName(), fieldElement.getTag(), fieldElement.getDefaultValue())); + if (!messageElementFields.isEmpty()) { + messageElementFields.forEach(fieldElement -> { + Field.Label label = fieldElement.getLabel(); + String labelStr = label != null ? label.name() : null; + messageDefinitionBuilder.addField( + labelStr, + fieldElement.getType(), + fieldElement.getName(), + fieldElement.getTag(), + fieldElement.getDefaultValue()); + }); } MessageDefinition messageDefinition = messageDefinitionBuilder.build(); schemaBuilder.addMessageDefinition(messageDefinition); @@ -129,16 +137,21 @@ public class MqttProtoDeviceProfileTransportConfiguration extends MqttDeviceProf DynamicMessage.Builder builder = dynamicSchema.newMessageBuilder(lastMsg.getName()); return builder.getDescriptorForType(); } catch (Descriptors.DescriptorValidationException e) { - throw new RuntimeException(e); + log.error("Failed to create dynamic schema due to: ", e); + return null; } } else { - throw new RuntimeException("Failed to get Message Descriptor! Message types is empty for " + schemaName + " schema!"); + log.error("Failed to get Message Descriptor! Message types is empty for {} schema!", schemaName); + return null; } } - private ProtoFileElement getTransportProtoSchema(String protoSchema) { return new ProtoParser(LOCATION, protoSchema.toCharArray()).readProtoFile(); } + private boolean isEmptyStr(String str) { + return str == null || "".equals(str); + } + } 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 3b3f80e89f..fab0579bb2 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 @@ -262,24 +262,24 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { try { - MqttTransportAdaptor adaptor = deviceSessionCtx.getPayloadAdaptor(); + MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor(); if (deviceSessionCtx.isDeviceTelemetryTopic(topicName)) { - TransportProtos.PostTelemetryMsg postTelemetryMsg = adaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); + TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { - TransportProtos.PostAttributeMsg postAttributeMsg = adaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); + TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { - TransportProtos.GetAttributeRequestMsg getAttributeMsg = adaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg); + TransportProtos.GetAttributeRequestMsg getAttributeMsg = payloadAdaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_RPC_RESPONSE_TOPIC)) { - TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = adaptor.convertToDeviceRpcResponse(deviceSessionCtx, mqttMsg); + TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = payloadAdaptor.convertToDeviceRpcResponse(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_RPC_REQUESTS_TOPIC)) { - TransportProtos.ToServerRpcRequestMsg rpcRequestMsg = adaptor.convertToServerRpcRequest(deviceSessionCtx, mqttMsg); + TransportProtos.ToServerRpcRequestMsg rpcRequestMsg = payloadAdaptor.convertToServerRpcRequest(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequestMsg, getPubAckCallback(ctx, msgId, rpcRequestMsg)); } else if (topicName.equals(MqttTopics.DEVICE_CLAIM_TOPIC)) { - TransportProtos.ClaimDeviceMsg claimDeviceMsg = adaptor.convertToClaimDevice(deviceSessionCtx, mqttMsg); + TransportProtos.ClaimDeviceMsg claimDeviceMsg = payloadAdaptor.convertToClaimDevice(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), claimDeviceMsg, getPubAckCallback(ctx, msgId, claimDeviceMsg)); } else { transportService.reportActivity(deviceSessionCtx.getSessionInfo()); 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 532115a765..aaa8983a0c 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 @@ -38,7 +38,6 @@ import org.thingsboard.server.common.transport.adaptor.ProtoConverter; import org.thingsboard.server.gen.transport.TransportApiProtos; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; -import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg; import org.thingsboard.server.transport.mqtt.session.MqttDeviceAwareSessionContext; import java.util.Optional; @@ -122,7 +121,6 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { @Override public TransportProtos.ProvisionDeviceRequestMsg convertToProvisionRequestMsg(MqttDeviceAwareSessionContext ctx, MqttPublishMessage mqttMsg) throws AdaptorException { byte[] bytes = toBytes(mqttMsg.payload()); - String topicName = mqttMsg.variableHeader().topicName(); try { return ProtoConverter.convertToProvisionRequestMsg(bytes); } catch (InvalidProtocolBufferException ex) { 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 f595df7721..6ecea5a77c 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 @@ -121,7 +121,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { payloadType = mqttConfig.getTransportPayloadType(); telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); attributesTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); - if (payloadType.equals(TransportPayloadType.PROTOBUF) && mqttConfig instanceof MqttProtoDeviceProfileTransportConfiguration) { + if (mqttConfig instanceof MqttProtoDeviceProfileTransportConfiguration) { updateDynamicMessageDescriptors(mqttConfig); } } else {