@ -22,6 +22,7 @@ import com.squareup.wire.schema.internal.parser.ProtoFileElement;
import lombok.extern.slf4j.Slf4j ;
import org.eclipse.paho.client.mqttv3.MqttAsyncClient ;
import org.junit.After ;
import org.junit.Ignore ;
import org.junit.Test ;
import org.thingsboard.server.common.data.Device ;
import org.thingsboard.server.common.data.DeviceProfileProvisionType ;
@ -35,6 +36,7 @@ import org.thingsboard.server.gen.transport.TransportApiProtos;
import org.thingsboard.server.gen.transport.TransportProtos ;
import java.util.Arrays ;
import java.util.Collections ;
import java.util.List ;
import static org.junit.Assert.assertNotNull ;
@ -51,9 +53,8 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac
}
@Test
public void testPushMqtt Telemetry ( ) throws Exception {
public void testPushTelemetry ( ) throws Exception {
super . processBeforeTest ( "Test Post Telemetry device proto payload" , "Test Post Telemetry gateway proto payload" , TransportPayloadType . PROTOBUF , POST_DATA_TELEMETRY_TOPIC , null ) ;
List < String > expectedKeys = Arrays . asList ( "key1" , "key2" , "key3" , "key4" , "key5" ) ;
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile . getProfileData ( ) . getTransportConfiguration ( ) ;
assertTrue ( transportConfiguration instanceof MqttDeviceProfileTransportConfiguration ) ;
MqttDeviceProfileTransportConfiguration mqttTransportConfiguration = ( MqttDeviceProfileTransportConfiguration ) transportConfiguration ;
@ -89,38 +90,37 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key4" ) , 4 )
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key5" ) , jsonObject )
. build ( ) ;
processTelemetryTest ( POST_DATA_TELEMETRY_TOPIC , expectedKeys , postTelemetryMsg . toByteArray ( ) , false ) ;
processTelemetryTest ( POST_DATA_TELEMETRY_TOPIC , Arrays . asList ( "key1" , "key2" , "key3" , "key4" , "key5" ) , postTelemetryMsg . toByteArray ( ) , false , false ) ;
}
@Test
public void testPushMqtt TelemetryWithTs ( ) throws Exception {
public void testPushTelemetryWithTs ( ) throws Exception {
String schemaStr = "syntax =\"proto3\";\n" +
"\n" +
"package test;\n" +
"\n" +
"message PostTelemetry {\n" +
" int64 ts = 1;\n" +
" optional int64 ts = 1;\n" +
" Values values = 2;\n" +
" \n" +
" message Values {\n" +
" string key1 = 3;\n" +
" bool key2 = 4;\n" +
" double key3 = 5;\n" +
" int32 key4 = 6;\n" +
" optional string key1 = 3;\n" +
" optional bool key2 = 4;\n" +
" optional double key3 = 5;\n" +
" optional int32 key4 = 6;\n" +
" JsonObject key5 = 7;\n" +
" }\n" +
" \n" +
" message JsonObject {\n" +
" int32 someNumber = 8;\n" +
" optional int32 someNumber = 8;\n" +
" repeated int32 someArray = 9;\n" +
" NestedJsonObject someNestedObject = 10;\n" +
" message NestedJsonObject {\n" +
" string key = 11;\n" +
" optional string key = 11;\n" +
" }\n" +
" }\n" +
"}" ;
super . processBeforeTest ( "Test Post Telemetry device proto payload" , "Test Post Telemetry gateway proto payload" , TransportPayloadType . PROTOBUF , POST_DATA_TELEMETRY_TOPIC , null , schemaStr , null , null , null , null , null , DeviceProfileProvisionType . DISABLED ) ;
List < String > expectedKeys = Arrays . asList ( "key1" , "key2" , "key3" , "key4" , "key5" ) ;
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile . getProfileData ( ) . getTransportConfiguration ( ) ;
assertTrue ( transportConfiguration instanceof MqttDeviceProfileTransportConfiguration ) ;
MqttDeviceProfileTransportConfiguration mqttTransportConfiguration = ( MqttDeviceProfileTransportConfiguration ) transportConfiguration ;
@ -167,11 +167,124 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "values" ) , valuesMsg )
. build ( ) ;
processTelemetryTest ( POST_DATA_TELEMETRY_TOPIC , expectedKeys , postTelemetryMsg . toByteArray ( ) , true ) ;
processTelemetryTest ( POST_DATA_TELEMETRY_TOPIC , Arrays . asList ( "key1" , "key2" , "key3" , "key4" , "key5" ) , postTelemetryMsg . toByteArray ( ) , true , false ) ;
}
@Test
public void testPushTelemetryWithExplicitPresenceProtoKeys ( ) throws Exception {
super . processBeforeTest ( "Test Post Telemetry device proto payload" , "Test Post Telemetry gateway proto payload" , TransportPayloadType . PROTOBUF , POST_DATA_TELEMETRY_TOPIC , null ) ;
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile . getProfileData ( ) . getTransportConfiguration ( ) ;
assertTrue ( transportConfiguration instanceof MqttDeviceProfileTransportConfiguration ) ;
MqttDeviceProfileTransportConfiguration mqttTransportConfiguration = ( MqttDeviceProfileTransportConfiguration ) transportConfiguration ;
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttTransportConfiguration . getTransportPayloadTypeConfiguration ( ) ;
assertTrue ( transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration ) ;
ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = ( ProtoTransportPayloadConfiguration ) transportPayloadTypeConfiguration ;
ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration . getTransportProtoSchema ( DEVICE_TELEMETRY_PROTO_SCHEMA ) ;
DynamicSchema telemetrySchema = protoTransportPayloadConfiguration . getDynamicSchema ( transportProtoSchema , "telemetrySchema" ) ;
DynamicMessage . Builder nestedJsonObjectBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry.JsonObject.NestedJsonObject" ) ;
Descriptors . Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder . getDescriptorForType ( ) ;
assertNotNull ( nestedJsonObjectBuilderDescriptor ) ;
DynamicMessage nestedJsonObject = nestedJsonObjectBuilder . setField ( nestedJsonObjectBuilderDescriptor . findFieldByName ( "key" ) , "value" ) . build ( ) ;
DynamicMessage . Builder jsonObjectBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry.JsonObject" ) ;
Descriptors . Descriptor jsonObjectBuilderDescriptor = jsonObjectBuilder . getDescriptorForType ( ) ;
assertNotNull ( jsonObjectBuilderDescriptor ) ;
DynamicMessage jsonObject = jsonObjectBuilder
. addRepeatedField ( jsonObjectBuilderDescriptor . findFieldByName ( "someArray" ) , 1 )
. addRepeatedField ( jsonObjectBuilderDescriptor . findFieldByName ( "someArray" ) , 2 )
. addRepeatedField ( jsonObjectBuilderDescriptor . findFieldByName ( "someArray" ) , 3 )
. setField ( jsonObjectBuilderDescriptor . findFieldByName ( "someNestedObject" ) , nestedJsonObject )
. build ( ) ;
DynamicMessage . Builder postTelemetryBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry" ) ;
Descriptors . Descriptor postTelemetryMsgDescriptor = postTelemetryBuilder . getDescriptorForType ( ) ;
assertNotNull ( postTelemetryMsgDescriptor ) ;
DynamicMessage postTelemetryMsg = postTelemetryBuilder
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key1" ) , "" )
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key2" ) , false )
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key3" ) , 0 . 0 )
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key4" ) , 0 )
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "key5" ) , jsonObject )
. build ( ) ;
processTelemetryTest ( POST_DATA_TELEMETRY_TOPIC , Arrays . asList ( "key1" , "key2" , "key3" , "key4" , "key5" ) , postTelemetryMsg . toByteArray ( ) , false , true ) ;
}
@Test
public void testPushTelemetryWithTsAndNoPresenceFields ( ) throws Exception {
String schemaStr = "syntax =\"proto3\";\n" +
"\n" +
"package test;\n" +
"\n" +
"message PostTelemetry {\n" +
" optional int64 ts = 1;\n" +
" Values values = 2;\n" +
" \n" +
" message Values {\n" +
" string key1 = 3;\n" +
" bool key2 = 4;\n" +
" double key3 = 5;\n" +
" int32 key4 = 6;\n" +
" JsonObject key5 = 7;\n" +
" }\n" +
" \n" +
" message JsonObject {\n" +
" optional int32 someNumber = 8;\n" +
" repeated int32 someArray = 9;\n" +
" NestedJsonObject someNestedObject = 10;\n" +
" message NestedJsonObject {\n" +
" optional string key = 11;\n" +
" }\n" +
" }\n" +
"}" ;
super . processBeforeTest ( "Test Post Telemetry device proto payload" , "Test Post Telemetry gateway proto payload" , TransportPayloadType . PROTOBUF , POST_DATA_TELEMETRY_TOPIC , null , schemaStr , null , null , null , null , null , DeviceProfileProvisionType . DISABLED ) ;
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile . getProfileData ( ) . getTransportConfiguration ( ) ;
assertTrue ( transportConfiguration instanceof MqttDeviceProfileTransportConfiguration ) ;
MqttDeviceProfileTransportConfiguration mqttTransportConfiguration = ( MqttDeviceProfileTransportConfiguration ) transportConfiguration ;
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttTransportConfiguration . getTransportPayloadTypeConfiguration ( ) ;
assertTrue ( transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration ) ;
ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = ( ProtoTransportPayloadConfiguration ) transportPayloadTypeConfiguration ;
ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration . getTransportProtoSchema ( schemaStr ) ;
DynamicSchema telemetrySchema = protoTransportPayloadConfiguration . getDynamicSchema ( transportProtoSchema , "telemetrySchema" ) ;
DynamicMessage . Builder nestedJsonObjectBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry.JsonObject.NestedJsonObject" ) ;
Descriptors . Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder . getDescriptorForType ( ) ;
assertNotNull ( nestedJsonObjectBuilderDescriptor ) ;
DynamicMessage nestedJsonObject = nestedJsonObjectBuilder . setField ( nestedJsonObjectBuilderDescriptor . findFieldByName ( "key" ) , "value" ) . build ( ) ;
DynamicMessage . Builder jsonObjectBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry.JsonObject" ) ;
Descriptors . Descriptor jsonObjectBuilderDescriptor = jsonObjectBuilder . getDescriptorForType ( ) ;
assertNotNull ( jsonObjectBuilderDescriptor ) ;
DynamicMessage jsonObject = jsonObjectBuilder
. addRepeatedField ( jsonObjectBuilderDescriptor . findFieldByName ( "someArray" ) , 1 )
. addRepeatedField ( jsonObjectBuilderDescriptor . findFieldByName ( "someArray" ) , 2 )
. addRepeatedField ( jsonObjectBuilderDescriptor . findFieldByName ( "someArray" ) , 3 )
. setField ( jsonObjectBuilderDescriptor . findFieldByName ( "someNestedObject" ) , nestedJsonObject )
. build ( ) ;
DynamicMessage . Builder valuesBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry.Values" ) ;
Descriptors . Descriptor valuesDescriptor = valuesBuilder . getDescriptorForType ( ) ;
assertNotNull ( valuesDescriptor ) ;
DynamicMessage valuesMsg = valuesBuilder
. setField ( valuesDescriptor . findFieldByName ( "key4" ) , 0 )
. setField ( valuesDescriptor . findFieldByName ( "key5" ) , jsonObject )
. build ( ) ;
DynamicMessage . Builder postTelemetryBuilder = telemetrySchema . newMessageBuilder ( "PostTelemetry" ) ;
Descriptors . Descriptor postTelemetryMsgDescriptor = postTelemetryBuilder . getDescriptorForType ( ) ;
assertNotNull ( postTelemetryMsgDescriptor ) ;
DynamicMessage postTelemetryMsg = postTelemetryBuilder
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "ts" ) , 10000L )
. setField ( postTelemetryMsgDescriptor . findFieldByName ( "values" ) , valuesMsg )
. build ( ) ;
processTelemetryTest ( POST_DATA_TELEMETRY_TOPIC , Collections . singletonList ( "key5" ) , postTelemetryMsg . toByteArray ( ) , true , true ) ;
}
@Test
public void testPushMqttTelemetryGateway ( ) throws Exception {
public void testPushTelemetryGateway ( ) throws Exception {
super . processBeforeTest ( "Test Post Telemetry device proto payload" , "Test Post Telemetry gateway proto payload" , TransportPayloadType . PROTOBUF , null , null , null , null , null , null , null , null , DeviceProfileProvisionType . DISABLED ) ;
TransportApiProtos . GatewayTelemetryMsg . Builder gatewayTelemetryMsgProtoBuilder = TransportApiProtos . GatewayTelemetryMsg . newBuilder ( ) ;
List < String > expectedKeys = Arrays . asList ( "key1" , "key2" , "key3" , "key4" , "key5" ) ;