Browse Source

cleanup code

pull/3740/head
ShvaykaD 6 years ago
parent
commit
d11a1ba7e6
  1. 8
      application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java
  2. 4
      common/data/pom.xml
  3. 4
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java
  4. 37
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttProtoDeviceProfileTransportConfiguration.java
  5. 14
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  6. 2
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
  7. 2
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java

8
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));

4
common/data/pom.xml

@ -71,10 +71,6 @@
<artifactId>java-driver-core</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-web</artifactId>
</dependency>
<dependency>
<groupId>com.squareup.wire</groupId>
<artifactId>wire-schema</artifactId>

4
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;

37
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<TypeElement> 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<TypeElement> 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<EnumConstantElement> 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<FieldElement> 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);
}
}

14
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());

2
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) {

2
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 {

Loading…
Cancel
Save