diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java index b811497b9f..3fe80b1112 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java @@ -16,12 +16,6 @@ package org.thingsboard.server.controller; import com.fasterxml.jackson.core.type.TypeReference; -import com.github.os72.protobuf.dynamic.DynamicSchema; -import com.google.protobuf.Descriptors; -import com.google.protobuf.DynamicMessage; -import com.google.protobuf.InvalidProtocolBufferException; -import com.google.protobuf.util.JsonFormat; -import com.squareup.wire.schema.internal.parser.ProtoFileElement; import org.junit.After; import org.junit.Assert; import org.junit.Before; @@ -46,11 +40,9 @@ import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.audit.ActionType; -import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.JsonTransportPayloadConfiguration; import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration; -import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; @@ -62,13 +54,9 @@ import org.thingsboard.server.dao.exception.DataValidationException; import java.util.ArrayList; import java.util.Collections; import java.util.List; -import java.util.Set; import java.util.stream.Collectors; import static org.hamcrest.Matchers.containsString; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertTrue; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE; import static org.thingsboard.server.common.data.ota.OtaPackageType.SOFTWARE; @@ -749,137 +737,6 @@ public abstract class BaseDeviceProfileControllerTest extends AbstractController "}", "[Transport Configuration] invalid attributes proto schema provided! OneOf definition groups don't support!"); } - @Test - public void testSaveProtoDeviceProfileWithMessageNestedTypes() throws Exception { - String schema = "syntax = \"proto3\";\n" + - "\n" + - "package testnested;\n" + - "\n" + - "message Outer {\n" + - " message MiddleAA {\n" + - " message Inner {\n" + - " optional int64 ival = 1;\n" + - " optional bool booly = 2;\n" + - " }\n" + - " Inner inner = 1;\n" + - " }\n" + - " message MiddleBB {\n" + - " message Inner {\n" + - " optional int32 ival = 1;\n" + - " optional bool booly = 2;\n" + - " }\n" + - " Inner inner = 1;\n" + - " }\n" + - " MiddleAA middleAA = 1;\n" + - " MiddleBB middleBB = 2;\n" + - "}"; - DynamicSchema dynamicSchema = getDynamicSchema(schema); - assertNotNull(dynamicSchema); - Set messageTypes = dynamicSchema.getMessageTypes(); - assertEquals(5, messageTypes.size()); - assertTrue(messageTypes.contains("testnested.Outer")); - assertTrue(messageTypes.contains("testnested.Outer.MiddleAA")); - assertTrue(messageTypes.contains("testnested.Outer.MiddleAA.Inner")); - assertTrue(messageTypes.contains("testnested.Outer.MiddleBB")); - assertTrue(messageTypes.contains("testnested.Outer.MiddleBB.Inner")); - - DynamicMessage.Builder middleAAInnerMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleAA.Inner"); - Descriptors.Descriptor middleAAInnerMsgDescriptor = middleAAInnerMsgBuilder.getDescriptorForType(); - DynamicMessage middleAAInnerMsg = middleAAInnerMsgBuilder - .setField(middleAAInnerMsgDescriptor.findFieldByName("ival"), 1L) - .setField(middleAAInnerMsgDescriptor.findFieldByName("booly"), true) - .build(); - - DynamicMessage.Builder middleAAMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleAA"); - Descriptors.Descriptor middleAAMsgDescriptor = middleAAMsgBuilder.getDescriptorForType(); - DynamicMessage middleAAMsg = middleAAMsgBuilder - .setField(middleAAMsgDescriptor.findFieldByName("inner"), middleAAInnerMsg) - .build(); - - DynamicMessage.Builder middleBBInnerMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleAA.Inner"); - Descriptors.Descriptor middleBBInnerMsgDescriptor = middleBBInnerMsgBuilder.getDescriptorForType(); - DynamicMessage middleBBInnerMsg = middleBBInnerMsgBuilder - .setField(middleBBInnerMsgDescriptor.findFieldByName("ival"), 0L) - .setField(middleBBInnerMsgDescriptor.findFieldByName("booly"), false) - .build(); - - DynamicMessage.Builder middleBBMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleBB"); - Descriptors.Descriptor middleBBMsgDescriptor = middleBBMsgBuilder.getDescriptorForType(); - DynamicMessage middleBBMsg = middleBBMsgBuilder - .setField(middleBBMsgDescriptor.findFieldByName("inner"), middleBBInnerMsg) - .build(); - - - DynamicMessage.Builder outerMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer"); - Descriptors.Descriptor outerMsgBuilderDescriptor = outerMsgBuilder.getDescriptorForType(); - DynamicMessage outerMsg = outerMsgBuilder - .setField(outerMsgBuilderDescriptor.findFieldByName("middleAA"), middleAAMsg) - .setField(outerMsgBuilderDescriptor.findFieldByName("middleBB"), middleBBMsg) - .build(); - - assertEquals("{\n" + - " \"middleAA\": {\n" + - " \"inner\": {\n" + - " \"ival\": \"1\",\n" + - " \"booly\": true\n" + - " }\n" + - " },\n" + - " \"middleBB\": {\n" + - " \"inner\": {\n" + - " \"ival\": 0,\n" + - " \"booly\": false\n" + - " }\n" + - " }\n" + - "}", dynamicMsgToJson(outerMsgBuilderDescriptor, outerMsg.toByteArray())); - } - - @Test - public void testSaveProtoDeviceProfileWithMessageOneOfs() throws Exception { - String schema = "syntax = \"proto3\";\n" + - "\n" + - "package testoneofs;\n" + - "\n" + - "message SubMessage {\n" + - " repeated string name = 1;\n" + - "}\n" + - "\n" + - "message SampleMessage {\n" + - " optional int32 id = 1;\n" + - " oneof testOneOf {\n" + - " string name = 4;\n" + - " SubMessage subMessage = 9;\n" + - " }\n" + - "}"; - DynamicSchema dynamicSchema = getDynamicSchema(schema); - assertNotNull(dynamicSchema); - Set messageTypes = dynamicSchema.getMessageTypes(); - assertEquals(2, messageTypes.size()); - assertTrue(messageTypes.contains("testoneofs.SubMessage")); - assertTrue(messageTypes.contains("testoneofs.SampleMessage")); - - DynamicMessage.Builder sampleMsgBuilder = dynamicSchema.newMessageBuilder("testoneofs.SampleMessage"); - Descriptors.Descriptor sampleMsgDescriptor = sampleMsgBuilder.getDescriptorForType(); - assertNotNull(sampleMsgDescriptor); - - List fields = sampleMsgDescriptor.getFields(); - assertEquals(3, fields.size()); - DynamicMessage sampleMsg = sampleMsgBuilder - .setField(sampleMsgDescriptor.findFieldByName("name"), "Bob") - .build(); - assertEquals("{\n" + " \"name\": \"Bob\"\n" + "}", dynamicMsgToJson(sampleMsgDescriptor, sampleMsg.toByteArray())); - - DynamicMessage.Builder subMsgBuilder = dynamicSchema.newMessageBuilder("testoneofs.SubMessage"); - Descriptors.Descriptor subMsgDescriptor = subMsgBuilder.getDescriptorForType(); - DynamicMessage subMsg = subMsgBuilder - .addRepeatedField(subMsgDescriptor.findFieldByName("name"), "Alice") - .addRepeatedField(subMsgDescriptor.findFieldByName("name"), "John") - .build(); - - DynamicMessage sampleMsgWithOneOfSubMessage = sampleMsgBuilder.setField(sampleMsgDescriptor.findFieldByName("subMessage"), subMsg).build(); - assertEquals("{\n" + " \"subMessage\": {\n" + " \"name\": [\"Alice\", \"John\"]\n" + " }\n" + "}", - dynamicMsgToJson(sampleMsgDescriptor, sampleMsgWithOneOfSubMessage.toByteArray())); - } - @Test public void testSaveProtoDeviceProfileWithInvalidTelemetrySchemaTsField() throws Exception { testSaveDeviceProfileWithInvalidProtoSchema("syntax =\"proto3\";\n" + @@ -1127,23 +984,6 @@ public abstract class BaseDeviceProfileControllerTest extends AbstractController tenantAdmin.getId(), tenantAdmin.getEmail(), ActionType.ADDED, new DataValidationException(errorMsg)); } - private DynamicSchema getDynamicSchema(String schema) throws Exception { - DeviceProfile deviceProfile = testSaveDeviceProfileWithProtoPayloadType(schema); - DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); - assertTrue(transportConfiguration instanceof MqttDeviceProfileTransportConfiguration); - MqttDeviceProfileTransportConfiguration mqttDeviceProfileTransportConfiguration = (MqttDeviceProfileTransportConfiguration) transportConfiguration; - TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttDeviceProfileTransportConfiguration.getTransportPayloadTypeConfiguration(); - assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); - ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement protoFile = protoTransportPayloadConfiguration.getTransportProtoSchema(schema); - return protoTransportPayloadConfiguration.getDynamicSchema(protoFile, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); - } - - private String dynamicMsgToJson(Descriptors.Descriptor descriptor, byte[] payload) throws InvalidProtocolBufferException { - DynamicMessage dynamicMessage = DynamicMessage.parseFrom(descriptor, payload); - return JsonFormat.printer().includingDefaultValueFields().print(dynamicMessage); - } - @Test public void testDeleteDeviceProfileWithDeleteRelationsOk() throws Exception { DeviceProfileId deviceProfileId = savedDeviceProfile("DeviceProfile for Test WithRelationsOk").getId(); diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java index cd9ed239a5..7617f2b826 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java @@ -27,6 +27,7 @@ import org.eclipse.californium.core.CoapResponse; import org.eclipse.californium.core.coap.CoAP; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.CoapDeviceTypeConfiguration; import org.thingsboard.server.common.data.device.profile.DefaultCoapDeviceTypeConfiguration; @@ -124,7 +125,6 @@ public abstract class AbstractCoapAttributesIntegrationTest extends AbstractCoap } private byte[] getAttributesProtoPayloadBytes() { - DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); assertTrue(transportConfiguration instanceof CoapDeviceProfileTransportConfiguration); CoapDeviceProfileTransportConfiguration coapTransportConfiguration = (CoapDeviceProfileTransportConfiguration) transportConfiguration; @@ -134,8 +134,8 @@ public abstract class AbstractCoapAttributesIntegrationTest extends AbstractCoap TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(ATTRIBUTES_SCHEMA_STR); - DynamicSchema attributesSchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceAttributesProtoSchema()); + DynamicSchema attributesSchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); DynamicMessage.Builder nestedJsonObjectBuilder = attributesSchema.newMessageBuilder("PostAttributes.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java index 2d034d3b41..41001427a5 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java @@ -28,6 +28,7 @@ import org.eclipse.californium.core.CoapResponse; import org.eclipse.californium.core.coap.CoAP; import org.eclipse.californium.core.coap.MediaTypeRegistry; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.CoapDeviceTypeConfiguration; @@ -158,8 +159,8 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC protected void processOnLoadProtoResponse(CoapResponse response, CoapTestClient client, Integer observe, CountDownLatch latch) { ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = getProtoTransportPayloadConfiguration(); - ProtoFileElement rpcRequestProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(RPC_REQUEST_PROTO_SCHEMA); - DynamicSchema rpcRequestProtoSchema = protoTransportPayloadConfiguration.getDynamicSchema(rpcRequestProtoSchemaFile, ProtoTransportPayloadConfiguration.RPC_REQUEST_PROTO_SCHEMA); + ProtoFileElement rpcRequestProtoFileElement = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceRpcRequestProtoSchema()); + DynamicSchema rpcRequestProtoSchema = DynamicProtoUtils.getDynamicSchema(rpcRequestProtoFileElement, ProtoTransportPayloadConfiguration.RPC_REQUEST_PROTO_SCHEMA); byte[] requestPayload = response.getPayload(); DynamicMessage.Builder rpcRequestMsg = rpcRequestProtoSchema.newMessageBuilder("RpcRequestMsg"); @@ -168,8 +169,8 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC DynamicMessage dynamicMessage = DynamicMessage.parseFrom(rpcRequestMsgDescriptor, requestPayload); Descriptors.FieldDescriptor requestIdDescriptor = rpcRequestMsgDescriptor.findFieldByName("requestId"); int requestId = (int) dynamicMessage.getField(requestIdDescriptor); - ProtoFileElement rpcResponseProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_RPC_RESPONSE_PROTO_SCHEMA); - DynamicSchema rpcResponseProtoSchema = protoTransportPayloadConfiguration.getDynamicSchema(rpcResponseProtoSchemaFile, ProtoTransportPayloadConfiguration.RPC_RESPONSE_PROTO_SCHEMA); + ProtoFileElement rpcResponseProtoSchemaFile = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceRpcResponseProtoSchema()); + DynamicSchema rpcResponseProtoSchema = DynamicProtoUtils.getDynamicSchema(rpcResponseProtoSchemaFile, ProtoTransportPayloadConfiguration.RPC_RESPONSE_PROTO_SCHEMA); DynamicMessage.Builder rpcResponseBuilder = rpcResponseProtoSchema.newMessageBuilder("RpcResponseMsg"); Descriptors.Descriptor rpcResponseMsgDescriptor = rpcResponseBuilder.getDescriptorForType(); DynamicMessage rpcResponseMsg = rpcResponseBuilder diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/attributes/CoapAttributesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/attributes/CoapAttributesProtoIntegrationTest.java index 539f62f614..609a5fe672 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/attributes/CoapAttributesProtoIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/attributes/CoapAttributesProtoIntegrationTest.java @@ -23,6 +23,7 @@ import lombok.extern.slf4j.Slf4j; import org.junit.Before; import org.junit.Test; import org.thingsboard.server.common.data.CoapDeviceType; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.CoapDeviceTypeConfiguration; @@ -55,18 +56,7 @@ public class CoapAttributesProtoIntegrationTest extends CoapAttributesIntegratio @Test public void testPushAttributes() throws Exception { - DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); - assertTrue(transportConfiguration instanceof CoapDeviceProfileTransportConfiguration); - CoapDeviceProfileTransportConfiguration coapTransportConfiguration = (CoapDeviceProfileTransportConfiguration) transportConfiguration; - CoapDeviceTypeConfiguration coapDeviceTypeConfiguration = coapTransportConfiguration.getCoapDeviceTypeConfiguration(); - assertTrue(coapDeviceTypeConfiguration instanceof DefaultCoapDeviceTypeConfiguration); - DefaultCoapDeviceTypeConfiguration defaultCoapDeviceTypeConfiguration = (DefaultCoapDeviceTypeConfiguration) coapDeviceTypeConfiguration; - TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); - assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); - ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_ATTRIBUTES_PROTO_SCHEMA); - DynamicSchema attributesSchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchemaFile, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); - + DynamicSchema attributesSchema = getDynamicSchema(); DynamicMessage.Builder nestedJsonObjectBuilder = attributesSchema.newMessageBuilder("PostAttributes.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); assertNotNull(nestedJsonObjectBuilderDescriptor); @@ -98,18 +88,7 @@ public class CoapAttributesProtoIntegrationTest extends CoapAttributesIntegratio @Test public void testPushAttributesWithExplicitPresenceProtoKeys() throws Exception { - DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); - assertTrue(transportConfiguration instanceof CoapDeviceProfileTransportConfiguration); - CoapDeviceProfileTransportConfiguration coapTransportConfiguration = (CoapDeviceProfileTransportConfiguration) transportConfiguration; - CoapDeviceTypeConfiguration coapDeviceTypeConfiguration = coapTransportConfiguration.getCoapDeviceTypeConfiguration(); - assertTrue(coapDeviceTypeConfiguration instanceof DefaultCoapDeviceTypeConfiguration); - DefaultCoapDeviceTypeConfiguration defaultCoapDeviceTypeConfiguration = (DefaultCoapDeviceTypeConfiguration) coapDeviceTypeConfiguration; - TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); - assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); - ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_ATTRIBUTES_PROTO_SCHEMA); - DynamicSchema attributesSchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchemaFile, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); - + DynamicSchema attributesSchema = getDynamicSchema(); DynamicMessage.Builder nestedJsonObjectBuilder = attributesSchema.newMessageBuilder("PostAttributes.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); assertNotNull(nestedJsonObjectBuilderDescriptor); @@ -135,4 +114,19 @@ public class CoapAttributesProtoIntegrationTest extends CoapAttributesIntegratio processAttributesTest(Arrays.asList("key1", "key5"), postAttributesMsg.toByteArray(), true); } + private DynamicSchema getDynamicSchema() { + DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); + assertTrue(transportConfiguration instanceof CoapDeviceProfileTransportConfiguration); + CoapDeviceProfileTransportConfiguration coapTransportConfiguration = (CoapDeviceProfileTransportConfiguration) transportConfiguration; + CoapDeviceTypeConfiguration coapDeviceTypeConfiguration = coapTransportConfiguration.getCoapDeviceTypeConfiguration(); + assertTrue(coapDeviceTypeConfiguration instanceof DefaultCoapDeviceTypeConfiguration); + DefaultCoapDeviceTypeConfiguration defaultCoapDeviceTypeConfiguration = (DefaultCoapDeviceTypeConfiguration) coapDeviceTypeConfiguration; + TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); + assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); + ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; + String deviceAttributesProtoSchema = protoTransportPayloadConfiguration.getDeviceAttributesProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceAttributesProtoSchema); + return DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); + } + } diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/timeseries/AbstractCoapTimeseriesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/timeseries/AbstractCoapTimeseriesProtoIntegrationTest.java index fb6d8aa0da..536155f1cf 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/timeseries/AbstractCoapTimeseriesProtoIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/telemetry/timeseries/AbstractCoapTimeseriesProtoIntegrationTest.java @@ -23,6 +23,8 @@ import lombok.extern.slf4j.Slf4j; import org.junit.Before; import org.junit.Test; import org.thingsboard.server.common.data.CoapDeviceType; +import org.thingsboard.server.common.data.DeviceProfileProvisionType; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.CoapDeviceTypeConfiguration; @@ -63,8 +65,9 @@ public abstract class AbstractCoapTimeseriesProtoIntegrationTest extends Abstrac TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_TELEMETRY_PROTO_SCHEMA); - DynamicSchema telemetrySchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, "telemetrySchema"); + String deviceTelemetryProtoSchema = protoTransportPayloadConfiguration.getDeviceTelemetryProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceTelemetryProtoSchema); + DynamicSchema telemetrySchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.TELEMETRY_PROTO_SCHEMA); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); @@ -138,8 +141,9 @@ public abstract class AbstractCoapTimeseriesProtoIntegrationTest extends Abstrac TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(schemaStr); - DynamicSchema telemetrySchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, "telemetrySchema"); + String deviceTelemetryProtoSchema = protoTransportPayloadConfiguration.getDeviceTelemetryProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceTelemetryProtoSchema); + DynamicSchema telemetrySchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.TELEMETRY_PROTO_SCHEMA); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); @@ -198,8 +202,9 @@ public abstract class AbstractCoapTimeseriesProtoIntegrationTest extends Abstrac TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_TELEMETRY_PROTO_SCHEMA); - DynamicSchema telemetrySchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, "telemetrySchema"); + String deviceTelemetryProtoSchema = protoTransportPayloadConfiguration.getDeviceTelemetryProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceTelemetryProtoSchema); + DynamicSchema telemetrySchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.TELEMETRY_PROTO_SCHEMA); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); @@ -271,8 +276,9 @@ public abstract class AbstractCoapTimeseriesProtoIntegrationTest extends Abstrac TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(schemaStr); - DynamicSchema telemetrySchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, "telemetrySchema"); + String deviceTelemetryProtoSchema = protoTransportPayloadConfiguration.getDeviceTelemetryProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceTelemetryProtoSchema); + DynamicSchema telemetrySchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.TELEMETRY_PROTO_SCHEMA); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java index 6d030aca9f..56d9ef5e9e 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java @@ -25,6 +25,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.test.context.TestPropertySource; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; @@ -494,8 +495,8 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttTransportConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(ATTRIBUTES_SCHEMA_STR); - DynamicSchema attributesSchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceAttributesProtoSchema()); + DynamicSchema attributesSchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); DynamicMessage.Builder nestedJsonObjectBuilder = attributesSchema.newMessageBuilder("PostAttributes.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java index 067af192e4..2cdfb76645 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java @@ -29,6 +29,7 @@ import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TransportPayloadType; @@ -382,8 +383,8 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM protected byte[] processProtoMessageArrived(String requestTopic, MqttMessage mqttMessage) throws MqttException, InvalidProtocolBufferException { if (requestTopic.startsWith(BASE_DEVICE_API_TOPIC) || requestTopic.startsWith(BASE_DEVICE_API_TOPIC_V2)) { ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = getProtoTransportPayloadConfiguration(); - ProtoFileElement rpcRequestProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(RPC_REQUEST_PROTO_SCHEMA); - DynamicSchema rpcRequestProtoSchema = protoTransportPayloadConfiguration.getDynamicSchema(rpcRequestProtoSchemaFile, ProtoTransportPayloadConfiguration.RPC_REQUEST_PROTO_SCHEMA); + ProtoFileElement rpcRequestProtoFileElement = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceRpcRequestProtoSchema()); + DynamicSchema rpcRequestProtoSchema = DynamicProtoUtils.getDynamicSchema(rpcRequestProtoFileElement, ProtoTransportPayloadConfiguration.RPC_REQUEST_PROTO_SCHEMA); byte[] requestPayload = mqttMessage.getPayload(); DynamicMessage.Builder rpcRequestMsg = rpcRequestProtoSchema.newMessageBuilder("RpcRequestMsg"); @@ -395,8 +396,8 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM for (Descriptors.FieldDescriptor fieldDescriptor: fields) { assertTrue(dynamicMessage.hasField(fieldDescriptor)); } - ProtoFileElement transportProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_RPC_RESPONSE_PROTO_SCHEMA); - DynamicSchema rpcResponseProtoSchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchemaFile, ProtoTransportPayloadConfiguration.RPC_RESPONSE_PROTO_SCHEMA); + ProtoFileElement rpcResponseProtoFileElement = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceRpcResponseProtoSchema()); + DynamicSchema rpcResponseProtoSchema = DynamicProtoUtils.getDynamicSchema(rpcResponseProtoFileElement, ProtoTransportPayloadConfiguration.RPC_RESPONSE_PROTO_SCHEMA); DynamicMessage.Builder rpcResponseBuilder = rpcResponseProtoSchema.newMessageBuilder("RpcResponseMsg"); Descriptors.Descriptor rpcResponseMsgDescriptor = rpcResponseBuilder.getDescriptorForType(); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java index 4904ed7a5d..7b2fa1143d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java @@ -22,6 +22,7 @@ import com.squareup.wire.schema.internal.parser.ProtoFileElement; import lombok.extern.slf4j.Slf4j; import org.junit.Before; import org.junit.Test; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; @@ -176,8 +177,9 @@ public class MqttAttributesProtoIntegrationTest extends MqttAttributesIntegratio TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttTransportConfiguration.getTransportPayloadTypeConfiguration(); assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; - ProtoFileElement transportProtoSchemaFile = protoTransportPayloadConfiguration.getTransportProtoSchema(DEVICE_ATTRIBUTES_PROTO_SCHEMA); - return protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchemaFile, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); + String deviceAttributesProtoSchema = protoTransportPayloadConfiguration.getDeviceAttributesProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceAttributesProtoSchema); + return DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); } private DynamicMessage getDefaultDynamicMessage() { diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java index 95de5f0069..9c92e3f4f0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java @@ -23,6 +23,7 @@ import lombok.extern.slf4j.Slf4j; import org.junit.Before; import org.junit.Test; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; @@ -117,7 +118,7 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac .telemetryProtoSchema(schemaStr) .build(); processBeforeTest(configProperties); - DynamicSchema telemetrySchema = getDynamicSchema(schemaStr); + DynamicSchema telemetrySchema = getDynamicSchema(); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); @@ -167,7 +168,7 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac .telemetryTopicFilter(POST_DATA_TELEMETRY_TOPIC) .build(); processBeforeTest(configProperties); - DynamicSchema telemetrySchema = getDynamicSchema(DEVICE_TELEMETRY_PROTO_SCHEMA); + DynamicSchema telemetrySchema = getDynamicSchema(); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); @@ -230,7 +231,7 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac .telemetryProtoSchema(schemaStr) .build(); processBeforeTest(configProperties); - DynamicSchema telemetrySchema = getDynamicSchema(schemaStr); + DynamicSchema telemetrySchema = getDynamicSchema(); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); @@ -458,19 +459,20 @@ public abstract class AbstractMqttTimeseriesProtoIntegrationTest extends Abstrac assertFalse(callback.isPubAckReceived()); } - private DynamicSchema getDynamicSchema(String deviceTelemetryProtoSchema) { + private DynamicSchema getDynamicSchema() { 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(deviceTelemetryProtoSchema); - return protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, "telemetrySchema"); + String deviceTelemetryProtoSchema = protoTransportPayloadConfiguration.getDeviceTelemetryProtoSchema(); + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(deviceTelemetryProtoSchema); + return DynamicProtoUtils.getDynamicSchema(protoFileElement, ProtoTransportPayloadConfiguration.TELEMETRY_PROTO_SCHEMA); } private DynamicMessage getDefaultDynamicMessage() { - DynamicSchema telemetrySchema = getDynamicSchema(DEVICE_TELEMETRY_PROTO_SCHEMA); + DynamicSchema telemetrySchema = getDynamicSchema(); DynamicMessage.Builder nestedJsonObjectBuilder = telemetrySchema.newMessageBuilder("PostTelemetry.JsonObject.NestedJsonObject"); Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); diff --git a/common/data/pom.xml b/common/data/pom.xml index decbba55aa..fd927c13d8 100644 --- a/common/data/pom.xml +++ b/common/data/pom.xml @@ -108,6 +108,10 @@ de.ruedigermoeller fst + + com.google.protobuf + protobuf-java-util + diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DynamicProtoUtils.java b/common/data/src/main/java/org/thingsboard/server/common/data/DynamicProtoUtils.java new file mode 100644 index 0000000000..8951222ebd --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DynamicProtoUtils.java @@ -0,0 +1,301 @@ +/** + * Copyright © 2016-2022 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.common.data; + +import com.github.os72.protobuf.dynamic.DynamicSchema; +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.google.protobuf.InvalidProtocolBufferException; +import com.google.protobuf.util.JsonFormat; +import com.squareup.wire.Syntax; +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; +import com.squareup.wire.schema.internal.parser.FieldElement; +import com.squareup.wire.schema.internal.parser.MessageElement; +import com.squareup.wire.schema.internal.parser.OneOfElement; +import com.squareup.wire.schema.internal.parser.ProtoFileElement; +import com.squareup.wire.schema.internal.parser.ProtoParser; +import com.squareup.wire.schema.internal.parser.TypeElement; +import lombok.extern.slf4j.Slf4j; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; + +@Slf4j +public class DynamicProtoUtils { + + public static final Location LOCATION = new Location("", "", -1, -1); + public static final String PROTO_3_SYNTAX = "proto3"; + + public static Descriptors.Descriptor getDescriptor(String protoSchema, String schemaName) { + try { + DynamicMessage.Builder builder = getDynamicMessageBuilder(protoSchema, schemaName); + return builder.getDescriptorForType(); + } catch (Exception e) { + log.warn("Failed to get Message Descriptor due to {}", e.getMessage()); + return null; + } + } + + public static DynamicMessage.Builder getDynamicMessageBuilder(String protoSchema, String schemaName) { + ProtoFileElement protoFileElement = getProtoFileElement(protoSchema); + DynamicSchema dynamicSchema = getDynamicSchema(protoFileElement, schemaName); + String lastMsgName = getMessageTypes(protoFileElement.getTypes()).stream() + .map(MessageElement::getName).reduce((previous, last) -> last).get(); + return dynamicSchema.newMessageBuilder(lastMsgName); + } + + public static DynamicSchema getDynamicSchema(ProtoFileElement protoFileElement, String schemaName) { + DynamicSchema.Builder schemaBuilder = DynamicSchema.newBuilder(); + schemaBuilder.setName(schemaName); + schemaBuilder.setSyntax(PROTO_3_SYNTAX); + schemaBuilder.setPackage(StringUtils.isNotEmpty(protoFileElement.getPackageName()) ? + protoFileElement.getPackageName() : schemaName.toLowerCase()); + List types = protoFileElement.getTypes(); + List messageTypes = getMessageTypes(types); + + if (!messageTypes.isEmpty()) { + List enumTypes = getEnumElements(types); + if (!enumTypes.isEmpty()) { + enumTypes.forEach(enumElement -> { + EnumDefinition enumDefinition = getEnumDefinition(enumElement); + schemaBuilder.addEnumDefinition(enumDefinition); + }); + } + List messageDefinitions = getMessageDefinitions(messageTypes); + messageDefinitions.forEach(schemaBuilder::addMessageDefinition); + try { + return schemaBuilder.build(); + } catch (Descriptors.DescriptorValidationException e) { + throw new RuntimeException("Failed to create dynamic schema due to: " + e.getMessage()); + } + } else { + throw new RuntimeException("Failed to get Dynamic Schema! Message types is empty for schema:" + schemaName); + } + } + + public static ProtoFileElement getProtoFileElement(String protoSchema) { + return new ProtoParser(LOCATION, protoSchema.toCharArray()).readProtoFile(); + } + + public static String dynamicMsgToJson(Descriptors.Descriptor descriptor, byte[] payload) throws InvalidProtocolBufferException { + DynamicMessage dynamicMessage = DynamicMessage.parseFrom(descriptor, payload); + return JsonFormat.printer().includingDefaultValueFields().print(dynamicMessage); + } + + public static DynamicMessage jsonToDynamicMessage(DynamicMessage.Builder builder, String payload) throws InvalidProtocolBufferException { + JsonFormat.parser().ignoringUnknownFields().merge(payload, builder); + return builder.build(); + } + + private static List getMessageTypes(List types) { + return types.stream() + .filter(typeElement -> typeElement instanceof MessageElement) + .map(typeElement -> (MessageElement) typeElement) + .collect(Collectors.toList()); + } + + private static List getEnumElements(List types) { + return types.stream() + .filter(typeElement -> typeElement instanceof EnumElement) + .map(typeElement -> (EnumElement) typeElement) + .collect(Collectors.toList()); + } + + private static List getMessageDefinitions(List messageElementsList) { + if (!messageElementsList.isEmpty()) { + List messageDefinitions = new ArrayList<>(); + messageElementsList.forEach(messageElement -> { + MessageDefinition.Builder messageDefinitionBuilder = MessageDefinition.newBuilder(messageElement.getName()); + + List nestedTypes = messageElement.getNestedTypes(); + if (!nestedTypes.isEmpty()) { + List nestedEnumTypes = getEnumElements(nestedTypes); + if (!nestedEnumTypes.isEmpty()) { + nestedEnumTypes.forEach(enumElement -> { + EnumDefinition nestedEnumDefinition = getEnumDefinition(enumElement); + messageDefinitionBuilder.addEnumDefinition(nestedEnumDefinition); + }); + } + List nestedMessageTypes = getMessageTypes(nestedTypes); + List nestedMessageDefinitions = getMessageDefinitions(nestedMessageTypes); + nestedMessageDefinitions.forEach(messageDefinitionBuilder::addMessageDefinition); + } + List messageElementFields = messageElement.getFields(); + List oneOfs = messageElement.getOneOfs(); + if (!oneOfs.isEmpty()) { + for (OneOfElement oneOfelement : oneOfs) { + MessageDefinition.OneofBuilder oneofBuilder = messageDefinitionBuilder.addOneof(oneOfelement.getName()); + addMessageFieldsToTheOneOfDefinition(oneOfelement.getFields(), oneofBuilder); + } + } + if (!messageElementFields.isEmpty()) { + addMessageFieldsToTheMessageDefinition(messageElementFields, messageDefinitionBuilder); + } + messageDefinitions.add(messageDefinitionBuilder.build()); + }); + return messageDefinitions; + } else { + return Collections.emptyList(); + } + } + + private static EnumDefinition getEnumDefinition(EnumElement enumElement) { + List enumElementTypeConstants = enumElement.getConstants(); + EnumDefinition.Builder enumDefinitionBuilder = EnumDefinition.newBuilder(enumElement.getName()); + if (!enumElementTypeConstants.isEmpty()) { + enumElementTypeConstants.forEach(constantElement -> enumDefinitionBuilder.addValue(constantElement.getName(), constantElement.getTag())); + } + return enumDefinitionBuilder.build(); + } + + + private static void addMessageFieldsToTheMessageDefinition(List messageElementFields, MessageDefinition.Builder messageDefinitionBuilder) { + messageElementFields.forEach(fieldElement -> { + String labelStr = null; + if (fieldElement.getLabel() != null) { + labelStr = fieldElement.getLabel().name().toLowerCase(); + } + messageDefinitionBuilder.addField( + labelStr, + fieldElement.getType(), + fieldElement.getName(), + fieldElement.getTag()); + }); + } + + private static void addMessageFieldsToTheOneOfDefinition(List oneOfsElementFields, MessageDefinition.OneofBuilder oneofBuilder) { + oneOfsElementFields.forEach(fieldElement -> oneofBuilder.addField( + fieldElement.getType(), + fieldElement.getName(), + fieldElement.getTag())); + oneofBuilder.msgDefBuilder(); + } + + // validation + + public static void validateProtoSchema(String schema, String schemaName, String exceptionPrefix) throws IllegalArgumentException { + ProtoParser schemaParser = new ProtoParser(LOCATION, schema.toCharArray()); + ProtoFileElement protoFileElement; + try { + protoFileElement = schemaParser.readProtoFile(); + } catch (Exception e) { + throw new IllegalArgumentException(exceptionPrefix + " failed to parse " + schemaName + " due to: " + e.getMessage()); + } + checkProtoFileSyntax(schemaName, protoFileElement); + checkProtoFileCommonSettings(schemaName, protoFileElement.getOptions().isEmpty(), " Schema options don't support!", exceptionPrefix); + checkProtoFileCommonSettings(schemaName, protoFileElement.getPublicImports().isEmpty(), " Schema public imports don't support!", exceptionPrefix); + checkProtoFileCommonSettings(schemaName, protoFileElement.getImports().isEmpty(), " Schema imports don't support!", exceptionPrefix); + checkProtoFileCommonSettings(schemaName, protoFileElement.getExtendDeclarations().isEmpty(), " Schema extend declarations don't support!", exceptionPrefix); + checkTypeElements(schemaName, protoFileElement, exceptionPrefix); + } + + private static void checkProtoFileSyntax(String schemaName, ProtoFileElement protoFileElement) { + if (protoFileElement.getSyntax() == null || !protoFileElement.getSyntax().equals(Syntax.PROTO_3)) { + throw new IllegalArgumentException("[Transport Configuration] invalid schema syntax: " + protoFileElement.getSyntax() + + " for " + schemaName + " provided! Only " + Syntax.PROTO_3 + " allowed!"); + } + } + + private static void checkProtoFileCommonSettings(String schemaName, boolean isEmptySettings, String invalidSettingsMessage, String exceptionPrefix) { + if (!isEmptySettings) { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + invalidSettingsMessage); + } + } + + private static void checkTypeElements(String schemaName, ProtoFileElement protoFileElement, String exceptionPrefix) { + List types = protoFileElement.getTypes(); + if (!types.isEmpty()) { + if (types.stream().noneMatch(typeElement -> typeElement instanceof MessageElement)) { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + " At least one Message definition should exists!"); + } else { + checkEnumElements(schemaName, getEnumElements(types), exceptionPrefix); + checkMessageElements(schemaName, getMessageTypes(types), exceptionPrefix); + } + } else { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + " Type elements is empty!"); + } + } + + private static void checkFieldElements(String schemaName, List fieldElements, String exceptionPrefix) { + if (!fieldElements.isEmpty()) { + boolean hasRequiredLabel = fieldElements.stream().anyMatch(fieldElement -> { + Field.Label label = fieldElement.getLabel(); + return label != null && label.equals(Field.Label.REQUIRED); + }); + if (hasRequiredLabel) { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + " Required labels are not supported!"); + } + boolean hasDefaultValue = fieldElements.stream().anyMatch(fieldElement -> fieldElement.getDefaultValue() != null); + if (hasDefaultValue) { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + " Default values are not supported!"); + } + } + } + + private static void checkEnumElements(String schemaName, List enumTypes, String exceptionPrefix) { + if (enumTypes.stream().anyMatch(enumElement -> !enumElement.getNestedTypes().isEmpty())) { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + " Nested types in Enum definitions are not supported!"); + } + if (enumTypes.stream().anyMatch(enumElement -> !enumElement.getOptions().isEmpty())) { + throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName, exceptionPrefix) + " Enum definitions options are not supported!"); + } + } + + private static void checkMessageElements(String schemaName, List messageElementsList, String exceptionPrefix) { + if (!messageElementsList.isEmpty()) { + messageElementsList.forEach(messageElement -> { + checkProtoFileCommonSettings(schemaName, messageElement.getGroups().isEmpty(), + " Message definition groups don't support!", exceptionPrefix); + checkProtoFileCommonSettings(schemaName, messageElement.getOptions().isEmpty(), + " Message definition options don't support!", exceptionPrefix); + checkProtoFileCommonSettings(schemaName, messageElement.getExtensions().isEmpty(), + " Message definition extensions don't support!", exceptionPrefix); + checkProtoFileCommonSettings(schemaName, messageElement.getReserveds().isEmpty(), + " Message definition reserved elements don't support!", exceptionPrefix); + checkFieldElements(schemaName, messageElement.getFields(), exceptionPrefix); + List oneOfs = messageElement.getOneOfs(); + if (!oneOfs.isEmpty()) { + oneOfs.forEach(oneOfElement -> { + checkProtoFileCommonSettings(schemaName, oneOfElement.getGroups().isEmpty(), + " OneOf definition groups don't support!", exceptionPrefix); + checkFieldElements(schemaName, oneOfElement.getFields(), exceptionPrefix); + }); + } + List nestedTypes = messageElement.getNestedTypes(); + if (!nestedTypes.isEmpty()) { + List nestedEnumTypes = getEnumElements(nestedTypes); + if (!nestedEnumTypes.isEmpty()) { + checkEnumElements(schemaName, nestedEnumTypes, exceptionPrefix); + } + List nestedMessageTypes = getMessageTypes(nestedTypes); + checkMessageElements(schemaName, nestedMessageTypes, exceptionPrefix); + } + }); + } + } + + public static String invalidSchemaProvidedMessage(String schemaName, String exceptionPrefix) { + return exceptionPrefix + " invalid " + schemaName + " provided!"; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java index 409570887c..196069ff5f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/ProtoTransportPayloadConfiguration.java @@ -15,29 +15,15 @@ */ package org.thingsboard.server.common.data.device.profile; -import com.github.os72.protobuf.dynamic.DynamicSchema; -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.Location; -import com.squareup.wire.schema.internal.parser.EnumConstantElement; -import com.squareup.wire.schema.internal.parser.EnumElement; -import com.squareup.wire.schema.internal.parser.FieldElement; -import com.squareup.wire.schema.internal.parser.MessageElement; -import com.squareup.wire.schema.internal.parser.OneOfElement; -import com.squareup.wire.schema.internal.parser.ProtoFileElement; -import com.squareup.wire.schema.internal.parser.ProtoParser; -import com.squareup.wire.schema.internal.parser.TypeElement; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.DynamicProtoUtils; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TransportPayloadType; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.stream.Collectors; - @Slf4j @Data public class ProtoTransportPayloadConfiguration implements TransportPayloadTypeConfiguration { @@ -63,23 +49,23 @@ public class ProtoTransportPayloadConfiguration implements TransportPayloadTypeC } public Descriptors.Descriptor getTelemetryDynamicMessageDescriptor(String deviceTelemetryProtoSchema) { - return getDescriptor(deviceTelemetryProtoSchema, TELEMETRY_PROTO_SCHEMA); + return DynamicProtoUtils.getDescriptor(deviceTelemetryProtoSchema, TELEMETRY_PROTO_SCHEMA); } public Descriptors.Descriptor getAttributesDynamicMessageDescriptor(String deviceAttributesProtoSchema) { - return getDescriptor(deviceAttributesProtoSchema, ATTRIBUTES_PROTO_SCHEMA); + return DynamicProtoUtils.getDescriptor(deviceAttributesProtoSchema, ATTRIBUTES_PROTO_SCHEMA); } public Descriptors.Descriptor getRpcResponseDynamicMessageDescriptor(String deviceRpcResponseProtoSchema) { - return getDescriptor(deviceRpcResponseProtoSchema, RPC_RESPONSE_PROTO_SCHEMA); + return DynamicProtoUtils.getDescriptor(deviceRpcResponseProtoSchema, RPC_RESPONSE_PROTO_SCHEMA); } public DynamicMessage.Builder getRpcRequestDynamicMessageBuilder(String deviceRpcRequestProtoSchema) { - return getDynamicMessageBuilder(deviceRpcRequestProtoSchema, RPC_REQUEST_PROTO_SCHEMA); + return DynamicProtoUtils.getDynamicMessageBuilder(deviceRpcRequestProtoSchema, RPC_REQUEST_PROTO_SCHEMA); } public String getDeviceRpcResponseProtoSchema() { - if (!isEmptyStr(deviceRpcResponseProtoSchema)) { + if (StringUtils.isNotEmpty(deviceRpcResponseProtoSchema)) { return deviceRpcResponseProtoSchema; } else { return "syntax =\"proto3\";\n" + @@ -92,7 +78,7 @@ public class ProtoTransportPayloadConfiguration implements TransportPayloadTypeC } public String getDeviceRpcRequestProtoSchema() { - if (!isEmptyStr(deviceRpcRequestProtoSchema)) { + if (StringUtils.isNotEmpty(deviceRpcRequestProtoSchema)) { return deviceRpcRequestProtoSchema; } else { return "syntax =\"proto3\";\n" + @@ -106,143 +92,4 @@ public class ProtoTransportPayloadConfiguration implements TransportPayloadTypeC } } - private Descriptors.Descriptor getDescriptor(String protoSchema, String schemaName) { - try { - DynamicMessage.Builder builder = getDynamicMessageBuilder(protoSchema, schemaName); - return builder.getDescriptorForType(); - } catch (Exception e) { - log.warn("Failed to get Message Descriptor due to {}", e.getMessage()); - return null; - } - } - - public DynamicMessage.Builder getDynamicMessageBuilder(String protoSchema, String schemaName) { - ProtoFileElement protoFileElement = getTransportProtoSchema(protoSchema); - DynamicSchema dynamicSchema = getDynamicSchema(protoFileElement, schemaName); - String lastMsgName = getMessageTypes(protoFileElement.getTypes()).stream() - .map(MessageElement::getName).reduce((previous, last) -> last).get(); - return dynamicSchema.newMessageBuilder(lastMsgName); - } - - public DynamicSchema getDynamicSchema(ProtoFileElement protoFileElement, String schemaName) { - DynamicSchema.Builder schemaBuilder = DynamicSchema.newBuilder(); - schemaBuilder.setName(schemaName); - schemaBuilder.setSyntax(PROTO_3_SYNTAX); - schemaBuilder.setPackage(!isEmptyStr(protoFileElement.getPackageName()) ? - protoFileElement.getPackageName() : schemaName.toLowerCase()); - List types = protoFileElement.getTypes(); - List messageTypes = getMessageTypes(types); - - if (!messageTypes.isEmpty()) { - List enumTypes = getEnumElements(types); - if (!enumTypes.isEmpty()) { - enumTypes.forEach(enumElement -> { - EnumDefinition enumDefinition = getEnumDefinition(enumElement); - schemaBuilder.addEnumDefinition(enumDefinition); - }); - } - List messageDefinitions = getMessageDefinitions(messageTypes); - messageDefinitions.forEach(schemaBuilder::addMessageDefinition); - try { - return schemaBuilder.build(); - } catch (Descriptors.DescriptorValidationException e) { - throw new RuntimeException("Failed to create dynamic schema due to: " + e.getMessage()); - } - } else { - throw new RuntimeException("Failed to get Dynamic Schema! Message types is empty for schema:" + schemaName); - } - } - - public ProtoFileElement getTransportProtoSchema(String protoSchema) { - return new ProtoParser(LOCATION, protoSchema.toCharArray()).readProtoFile(); - } - - private List getMessageTypes(List types) { - return types.stream() - .filter(typeElement -> typeElement instanceof MessageElement) - .map(typeElement -> (MessageElement) typeElement) - .collect(Collectors.toList()); - } - - private List getEnumElements(List types) { - return types.stream() - .filter(typeElement -> typeElement instanceof EnumElement) - .map(typeElement -> (EnumElement) typeElement) - .collect(Collectors.toList()); - } - - private List getMessageDefinitions(List messageElementsList) { - if (!messageElementsList.isEmpty()) { - List messageDefinitions = new ArrayList<>(); - messageElementsList.forEach(messageElement -> { - MessageDefinition.Builder messageDefinitionBuilder = MessageDefinition.newBuilder(messageElement.getName()); - - List nestedTypes = messageElement.getNestedTypes(); - if (!nestedTypes.isEmpty()) { - List nestedEnumTypes = getEnumElements(nestedTypes); - if (!nestedEnumTypes.isEmpty()) { - nestedEnumTypes.forEach(enumElement -> { - EnumDefinition nestedEnumDefinition = getEnumDefinition(enumElement); - messageDefinitionBuilder.addEnumDefinition(nestedEnumDefinition); - }); - } - List nestedMessageTypes = getMessageTypes(nestedTypes); - List nestedMessageDefinitions = getMessageDefinitions(nestedMessageTypes); - nestedMessageDefinitions.forEach(messageDefinitionBuilder::addMessageDefinition); - } - List messageElementFields = messageElement.getFields(); - List oneOfs = messageElement.getOneOfs(); - if (!oneOfs.isEmpty()) { - for (OneOfElement oneOfelement : oneOfs) { - MessageDefinition.OneofBuilder oneofBuilder = messageDefinitionBuilder.addOneof(oneOfelement.getName()); - addMessageFieldsToTheOneOfDefinition(oneOfelement.getFields(), oneofBuilder); - } - } - if (!messageElementFields.isEmpty()) { - addMessageFieldsToTheMessageDefinition(messageElementFields, messageDefinitionBuilder); - } - messageDefinitions.add(messageDefinitionBuilder.build()); - }); - return messageDefinitions; - } else { - return Collections.emptyList(); - } - } - - private EnumDefinition getEnumDefinition(EnumElement enumElement) { - List enumElementTypeConstants = enumElement.getConstants(); - EnumDefinition.Builder enumDefinitionBuilder = EnumDefinition.newBuilder(enumElement.getName()); - if (!enumElementTypeConstants.isEmpty()) { - enumElementTypeConstants.forEach(constantElement -> enumDefinitionBuilder.addValue(constantElement.getName(), constantElement.getTag())); - } - return enumDefinitionBuilder.build(); - } - - - private void addMessageFieldsToTheMessageDefinition(List messageElementFields, MessageDefinition.Builder messageDefinitionBuilder) { - messageElementFields.forEach(fieldElement -> { - String labelStr = null; - if (fieldElement.getLabel() != null) { - labelStr = fieldElement.getLabel().name().toLowerCase(); - } - messageDefinitionBuilder.addField( - labelStr, - fieldElement.getType(), - fieldElement.getName(), - fieldElement.getTag()); - }); - } - - private void addMessageFieldsToTheOneOfDefinition(List oneOfsElementFields, MessageDefinition.OneofBuilder oneofBuilder) { - oneOfsElementFields.forEach(fieldElement -> oneofBuilder.addField( - fieldElement.getType(), - fieldElement.getName(), - fieldElement.getTag())); - oneofBuilder.msgDefBuilder(); - } - - private boolean isEmptyStr(String str) { - return str == null || "".equals(str); - } - } diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/DynamicProtoUtilsTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/DynamicProtoUtilsTest.java new file mode 100644 index 0000000000..9ae4651cc9 --- /dev/null +++ b/common/data/src/test/java/org/thingsboard/server/common/data/DynamicProtoUtilsTest.java @@ -0,0 +1,169 @@ +/** + * Copyright © 2016-2022 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.common.data; + +import com.github.os72.protobuf.dynamic.DynamicSchema; +import com.google.protobuf.Descriptors; +import com.google.protobuf.DynamicMessage; +import com.squareup.wire.schema.internal.parser.ProtoFileElement; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.junit.MockitoJUnitRunner; + +import java.util.List; +import java.util.Set; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; + +@RunWith(MockitoJUnitRunner.class) +public class DynamicProtoUtilsTest { + + @Test + public void testProtoSchemaWithMessageNestedTypes() throws Exception { + String schema = "syntax = \"proto3\";\n" + + "\n" + + "package testnested;\n" + + "\n" + + "message Outer {\n" + + " message MiddleAA {\n" + + " message Inner {\n" + + " optional int64 ival = 1;\n" + + " optional bool booly = 2;\n" + + " }\n" + + " Inner inner = 1;\n" + + " }\n" + + " message MiddleBB {\n" + + " message Inner {\n" + + " optional int32 ival = 1;\n" + + " optional bool booly = 2;\n" + + " }\n" + + " Inner inner = 1;\n" + + " }\n" + + " MiddleAA middleAA = 1;\n" + + " MiddleBB middleBB = 2;\n" + + "}"; + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(schema); + DynamicSchema dynamicSchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, "test schema with nested types"); + assertNotNull(dynamicSchema); + Set messageTypes = dynamicSchema.getMessageTypes(); + assertEquals(5, messageTypes.size()); + assertTrue(messageTypes.contains("testnested.Outer")); + assertTrue(messageTypes.contains("testnested.Outer.MiddleAA")); + assertTrue(messageTypes.contains("testnested.Outer.MiddleAA.Inner")); + assertTrue(messageTypes.contains("testnested.Outer.MiddleBB")); + assertTrue(messageTypes.contains("testnested.Outer.MiddleBB.Inner")); + + DynamicMessage.Builder middleAAInnerMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleAA.Inner"); + Descriptors.Descriptor middleAAInnerMsgDescriptor = middleAAInnerMsgBuilder.getDescriptorForType(); + DynamicMessage middleAAInnerMsg = middleAAInnerMsgBuilder + .setField(middleAAInnerMsgDescriptor.findFieldByName("ival"), 1L) + .setField(middleAAInnerMsgDescriptor.findFieldByName("booly"), true) + .build(); + + DynamicMessage.Builder middleAAMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleAA"); + Descriptors.Descriptor middleAAMsgDescriptor = middleAAMsgBuilder.getDescriptorForType(); + DynamicMessage middleAAMsg = middleAAMsgBuilder + .setField(middleAAMsgDescriptor.findFieldByName("inner"), middleAAInnerMsg) + .build(); + + DynamicMessage.Builder middleBBInnerMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleAA.Inner"); + Descriptors.Descriptor middleBBInnerMsgDescriptor = middleBBInnerMsgBuilder.getDescriptorForType(); + DynamicMessage middleBBInnerMsg = middleBBInnerMsgBuilder + .setField(middleBBInnerMsgDescriptor.findFieldByName("ival"), 0L) + .setField(middleBBInnerMsgDescriptor.findFieldByName("booly"), false) + .build(); + + DynamicMessage.Builder middleBBMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer.MiddleBB"); + Descriptors.Descriptor middleBBMsgDescriptor = middleBBMsgBuilder.getDescriptorForType(); + DynamicMessage middleBBMsg = middleBBMsgBuilder + .setField(middleBBMsgDescriptor.findFieldByName("inner"), middleBBInnerMsg) + .build(); + + + DynamicMessage.Builder outerMsgBuilder = dynamicSchema.newMessageBuilder("testnested.Outer"); + Descriptors.Descriptor outerMsgBuilderDescriptor = outerMsgBuilder.getDescriptorForType(); + DynamicMessage outerMsg = outerMsgBuilder + .setField(outerMsgBuilderDescriptor.findFieldByName("middleAA"), middleAAMsg) + .setField(outerMsgBuilderDescriptor.findFieldByName("middleBB"), middleBBMsg) + .build(); + + assertEquals("{\n" + + " \"middleAA\": {\n" + + " \"inner\": {\n" + + " \"ival\": \"1\",\n" + + " \"booly\": true\n" + + " }\n" + + " },\n" + + " \"middleBB\": {\n" + + " \"inner\": {\n" + + " \"ival\": 0,\n" + + " \"booly\": false\n" + + " }\n" + + " }\n" + + "}", DynamicProtoUtils.dynamicMsgToJson(outerMsgBuilderDescriptor, outerMsg.toByteArray())); + } + + @Test + public void testProtoSchemaWithMessageOneOfs() throws Exception { + String schema = "syntax = \"proto3\";\n" + + "\n" + + "package testoneofs;\n" + + "\n" + + "message SubMessage {\n" + + " repeated string name = 1;\n" + + "}\n" + + "\n" + + "message SampleMessage {\n" + + " optional int32 id = 1;\n" + + " oneof testOneOf {\n" + + " string name = 4;\n" + + " SubMessage subMessage = 9;\n" + + " }\n" + + "}"; + ProtoFileElement protoFileElement = DynamicProtoUtils.getProtoFileElement(schema); + DynamicSchema dynamicSchema = DynamicProtoUtils.getDynamicSchema(protoFileElement, "test schema with message oneOfs"); + assertNotNull(dynamicSchema); + Set messageTypes = dynamicSchema.getMessageTypes(); + assertEquals(2, messageTypes.size()); + assertTrue(messageTypes.contains("testoneofs.SubMessage")); + assertTrue(messageTypes.contains("testoneofs.SampleMessage")); + + DynamicMessage.Builder sampleMsgBuilder = dynamicSchema.newMessageBuilder("testoneofs.SampleMessage"); + Descriptors.Descriptor sampleMsgDescriptor = sampleMsgBuilder.getDescriptorForType(); + assertNotNull(sampleMsgDescriptor); + + List fields = sampleMsgDescriptor.getFields(); + assertEquals(3, fields.size()); + DynamicMessage sampleMsg = sampleMsgBuilder + .setField(sampleMsgDescriptor.findFieldByName("name"), "Bob") + .build(); + assertEquals("{\n" + " \"name\": \"Bob\"\n" + "}", DynamicProtoUtils.dynamicMsgToJson(sampleMsgDescriptor, sampleMsg.toByteArray())); + + DynamicMessage.Builder subMsgBuilder = dynamicSchema.newMessageBuilder("testoneofs.SubMessage"); + Descriptors.Descriptor subMsgDescriptor = subMsgBuilder.getDescriptorForType(); + DynamicMessage subMsg = subMsgBuilder + .addRepeatedField(subMsgDescriptor.findFieldByName("name"), "Alice") + .addRepeatedField(subMsgDescriptor.findFieldByName("name"), "John") + .build(); + + DynamicMessage sampleMsgWithOneOfSubMessage = sampleMsgBuilder.setField(sampleMsgDescriptor.findFieldByName("subMessage"), subMsg).build(); + assertEquals("{\n" + " \"subMessage\": {\n" + " \"name\": [\"Alice\", \"John\"]\n" + " }\n" + "}", + DynamicProtoUtils.dynamicMsgToJson(sampleMsgDescriptor, sampleMsgWithOneOfSubMessage.toByteArray())); + } + +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItem.java b/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItem.java index 7c7f300778..f17428f0bf 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItem.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItem.java @@ -20,10 +20,11 @@ import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.msg.gen.MsgProtos; +import java.io.Serializable; import java.util.UUID; @Data -public class TbMsgProcessingStackItem { +public class TbMsgProcessingStackItem implements Serializable { private final RuleChainId ruleChainId; private final RuleNodeId ruleNodeId; diff --git a/common/message/src/test/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItemTest.java b/common/message/src/test/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItemTest.java new file mode 100644 index 0000000000..8f8e59cd69 --- /dev/null +++ b/common/message/src/test/java/org/thingsboard/server/common/msg/TbMsgProcessingStackItemTest.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2022 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.common.msg; + +import org.junit.jupiter.api.Test; +import org.thingsboard.server.common.data.FSTUtils; +import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.RuleNodeId; + +import java.util.UUID; + +import static org.assertj.core.api.Assertions.assertThat; + +class TbMsgProcessingStackItemTest { + + @Test + void testSerialization() { + TbMsgProcessingStackItem item = new TbMsgProcessingStackItem(new RuleChainId(UUID.randomUUID()), new RuleNodeId(UUID.randomUUID())); + byte[] bytes = FSTUtils.encode(item); + TbMsgProcessingStackItem itemDecoded = FSTUtils.decode(bytes); + assertThat(item).isEqualTo(itemDecoded); + } + +} diff --git a/common/queue/pom.xml b/common/queue/pom.xml index c246ab358f..1ab33a3593 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -116,10 +116,6 @@ com.google.protobuf protobuf-java - - com.google.protobuf - protobuf-java-util - org.apache.curator curator-recipes diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/DefaultMvelInvokeService.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/DefaultMvelInvokeService.java index e72c24ecf4..393f935861 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/DefaultMvelInvokeService.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/DefaultMvelInvokeService.java @@ -24,7 +24,6 @@ import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.mvel2.ExecutionContext; import org.mvel2.MVEL; -import org.mvel2.ParserContext; import org.mvel2.SandboxedParserConfiguration; import org.mvel2.SandboxedParserContext; import org.mvel2.ScriptMemoryOverflowException; @@ -44,6 +43,7 @@ import org.thingsboard.server.common.stats.TbApiUsageStateClient; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.io.Serializable; +import java.util.Collections; import java.util.Map; import java.util.Optional; import java.util.UUID; @@ -115,7 +115,8 @@ public class DefaultMvelInvokeService extends AbstractScriptInvokeService implem executor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "mvel-executor")); try { // Special command to warm up MVEL engine - MVEL.compileExpression("var warmUp = {}; warmUp", new SandboxedParserContext(parserConfig)); + Serializable script = MVEL.compileExpression("var warmUp = {}; warmUp", new SandboxedParserContext(parserConfig)); + MVEL.executeTbExpression(script, new ExecutionContext(), Collections.emptyMap()); } catch (Exception e) { // do nothing } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java index 6fe0f29cfc..d7d685bd5a 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java @@ -27,6 +27,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.util.CollectionUtils; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.gen.transport.TransportApiProtos; import org.thingsboard.server.gen.transport.TransportProtos; @@ -184,8 +185,7 @@ public class ProtoConverter { try { JsonElement paramsElement = JSON_PARSER.parse(params); rpcRequestJson.add("params", paramsElement); - JsonFormat.parser().ignoringUnknownFields().merge(GSON.toJson(rpcRequestJson), rpcRequestDynamicMessageBuilder); - DynamicMessage dynamicRpcRequest = rpcRequestDynamicMessageBuilder.build(); + DynamicMessage dynamicRpcRequest = DynamicProtoUtils.jsonToDynamicMessage(rpcRequestDynamicMessageBuilder, GSON.toJson(rpcRequestJson)); return dynamicRpcRequest.toByteArray(); } catch (Exception e) { throw new AdaptorException("Failed to convert ToDeviceRpcRequestMsg to Dynamic Rpc request message due to: ", e); @@ -200,8 +200,7 @@ public class ProtoConverter { } public static String dynamicMsgToJson(byte[] bytes, Descriptors.Descriptor descriptor) throws InvalidProtocolBufferException { - DynamicMessage dynamicMessage = DynamicMessage.parseFrom(descriptor, bytes); - return JsonFormat.printer().includingDefaultValueFields().print(dynamicMessage); + return DynamicProtoUtils.dynamicMsgToJson(descriptor, bytes); } } diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index d703c53a7d..ec464caf7b 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -15,7 +15,9 @@ */ package org.thingsboard.common.util; +import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.json.JsonWriteFeature; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.MapperFeature; @@ -46,6 +48,10 @@ public class JacksonUtil { .configure(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS, true) .configure(MapperFeature.SORT_PROPERTIES_ALPHABETICALLY, true) .build(); + public static ObjectMapper ALLOW_UNQUOTED_FIELD_NAMES_MAPPER = JsonMapper.builder() + .configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false) + .configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true) + .build(); public static T convertValue(Object fromValue, Class toValueType) { try { @@ -127,18 +133,26 @@ public class JacksonUtil { } public static JsonNode toJsonNode(String value) { + return toJsonNode(value, OBJECT_MAPPER); + } + + public static JsonNode toJsonNode(String value, ObjectMapper mapper) { if (value == null || value.isEmpty()) { return null; } try { - return OBJECT_MAPPER.readTree(value); + return mapper.readTree(value); } catch (IOException e) { throw new IllegalArgumentException(e); } } public static ObjectNode newObjectNode() { - return OBJECT_MAPPER.createObjectNode(); + return newObjectNode(OBJECT_MAPPER); + } + + public static ObjectNode newObjectNode(ObjectMapper mapper) { + return mapper.createObjectNode(); } public static T clone(T value) { @@ -216,18 +230,27 @@ public class JacksonUtil { } public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) { + addKvEntry(entityNode, kvEntry, kvEntry.getKey()); + } + + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String key) { + addKvEntry(entityNode, kvEntry, key, OBJECT_MAPPER); + } + + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String key, ObjectMapper mapper) { if (kvEntry.getDataType() == DataType.BOOLEAN) { - kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); + kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.DOUBLE) { - kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); + kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.LONG) { - kvEntry.getLongValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); + kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.JSON) { if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(kvEntry.getKey(), JacksonUtil.toJsonNode(kvEntry.getJsonValue().get())); + entityNode.set(key, toJsonNode(kvEntry.getJsonValue().get(), mapper)); } } else { - entityNode.put(kvEntry.getKey(), kvEntry.getValueAsString()); + entityNode.put(key, kvEntry.getValueAsString()); } } + } diff --git a/common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java b/common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java new file mode 100644 index 0000000000..22cd1b3ff6 --- /dev/null +++ b/common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2022 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.common.util; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import org.junit.Assert; +import org.junit.Test; + +public class JacksonUtilTest { + + @Test + public void allow_unquoted_field_mapper_test() { + String data = "{data: 123}"; + JsonNode actualResult = JacksonUtil.toJsonNode(data, JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); // should be: {"data": 123} + ObjectNode expectedResult = JacksonUtil.newObjectNode(); + expectedResult.put("data", 123); // {"data": 123} + Assert.assertEquals(expectedResult, actualResult); + Assert.assertThrows(IllegalArgumentException.class, () -> JacksonUtil.toJsonNode(data)); // syntax exception due to missing quotes in the field name! + } + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java index 16c9bc4840..84eea112f3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java @@ -17,16 +17,6 @@ package org.thingsboard.server.dao.service.validator; import com.google.protobuf.Descriptors; import com.google.protobuf.DynamicMessage; -import com.squareup.wire.Syntax; -import com.squareup.wire.schema.Field; -import com.squareup.wire.schema.Location; -import com.squareup.wire.schema.internal.parser.EnumElement; -import com.squareup.wire.schema.internal.parser.FieldElement; -import com.squareup.wire.schema.internal.parser.MessageElement; -import com.squareup.wire.schema.internal.parser.OneOfElement; -import com.squareup.wire.schema.internal.parser.ProtoFileElement; -import com.squareup.wire.schema.internal.parser.ProtoParser; -import com.squareup.wire.schema.internal.parser.TypeElement; import org.eclipse.leshan.core.util.SecurityUtil; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; @@ -35,6 +25,7 @@ import org.springframework.util.CollectionUtils; import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileProvisionType; +import org.thingsboard.server.common.data.DynamicProtoUtils; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode; import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; @@ -67,16 +58,16 @@ import org.thingsboard.server.dao.tenant.TenantService; import java.util.HashSet; import java.util.List; import java.util.Set; -import java.util.stream.Collectors; @Component public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator { - private static final Location LOCATION = new Location("", "", -1, -1); private static final String ATTRIBUTES_PROTO_SCHEMA = "attributes proto schema"; private static final String TELEMETRY_PROTO_SCHEMA = "telemetry proto schema"; private static final String RPC_REQUEST_PROTO_SCHEMA = "rpc request proto schema"; private static final String RPC_RESPONSE_PROTO_SCHEMA = "rpc response proto schema"; + private static final String EXCEPTION_PREFIX = "[Transport Configuration]"; + @Autowired private DeviceProfileDao deviceProfileDao; @Autowired @@ -94,10 +85,6 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator types = protoFileElement.getTypes(); - if (!types.isEmpty()) { - if (types.stream().noneMatch(typeElement -> typeElement instanceof MessageElement)) { - throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName) + " At least one Message definition should exists!"); - } else { - checkEnumElements(schemaName, getEnumElements(types)); - checkMessageElements(schemaName, getMessageTypes(types)); - } - } else { - throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName) + " Type elements is empty!"); - } - } - - private void checkFieldElements(String schemaName, List fieldElements) { - if (!fieldElements.isEmpty()) { - boolean hasRequiredLabel = fieldElements.stream().anyMatch(fieldElement -> { - Field.Label label = fieldElement.getLabel(); - return label != null && label.equals(Field.Label.REQUIRED); - }); - if (hasRequiredLabel) { - throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName) + " Required labels are not supported!"); - } - boolean hasDefaultValue = fieldElements.stream().anyMatch(fieldElement -> fieldElement.getDefaultValue() != null); - if (hasDefaultValue) { - throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName) + " Default values are not supported!"); - } - } - } - - private void checkEnumElements(String schemaName, List enumTypes) { - if (enumTypes.stream().anyMatch(enumElement -> !enumElement.getNestedTypes().isEmpty())) { - throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName) + " Nested types in Enum definitions are not supported!"); - } - if (enumTypes.stream().anyMatch(enumElement -> !enumElement.getOptions().isEmpty())) { - throw new IllegalArgumentException(invalidSchemaProvidedMessage(schemaName) + " Enum definitions options are not supported!"); - } - } - - private void checkMessageElements(String schemaName, List messageElementsList) { - if (!messageElementsList.isEmpty()) { - messageElementsList.forEach(messageElement -> { - checkProtoFileCommonSettings(schemaName, messageElement.getGroups().isEmpty(), - " Message definition groups don't support!"); - checkProtoFileCommonSettings(schemaName, messageElement.getOptions().isEmpty(), - " Message definition options don't support!"); - checkProtoFileCommonSettings(schemaName, messageElement.getExtensions().isEmpty(), - " Message definition extensions don't support!"); - checkProtoFileCommonSettings(schemaName, messageElement.getReserveds().isEmpty(), - " Message definition reserved elements don't support!"); - checkFieldElements(schemaName, messageElement.getFields()); - List oneOfs = messageElement.getOneOfs(); - if (!oneOfs.isEmpty()) { - oneOfs.forEach(oneOfElement -> { - checkProtoFileCommonSettings(schemaName, oneOfElement.getGroups().isEmpty(), - " OneOf definition groups don't support!"); - checkFieldElements(schemaName, oneOfElement.getFields()); - }); - } - List nestedTypes = messageElement.getNestedTypes(); - if (!nestedTypes.isEmpty()) { - List nestedEnumTypes = getEnumElements(nestedTypes); - if (!nestedEnumTypes.isEmpty()) { - checkEnumElements(schemaName, nestedEnumTypes); - } - List nestedMessageTypes = getMessageTypes(nestedTypes); - checkMessageElements(schemaName, nestedMessageTypes); - } - }); - } - } - - private List getMessageTypes(List types) { - return types.stream() - .filter(typeElement -> typeElement instanceof MessageElement) - .map(typeElement -> (MessageElement) typeElement) - .collect(Collectors.toList()); - } - - private List getEnumElements(List types) { - return types.stream() - .filter(typeElement -> typeElement instanceof EnumElement) - .map(typeElement -> (EnumElement) typeElement) - .collect(Collectors.toList()); - } private void validateTelemetryDynamicMessageFields(ProtoTransportPayloadConfiguration protoTransportPayloadTypeConfiguration) { String deviceTelemetryProtoSchema = protoTransportPayloadTypeConfiguration.getDeviceTelemetryProtoSchema(); Descriptors.Descriptor telemetryDynamicMessageDescriptor = protoTransportPayloadTypeConfiguration.getTelemetryDynamicMessageDescriptor(deviceTelemetryProtoSchema); if (telemetryDynamicMessageDescriptor == null) { - throw new DataValidationException(invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA) + " Failed to get telemetryDynamicMessageDescriptor!"); + throw new DataValidationException(DynamicProtoUtils.invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA, EXCEPTION_PREFIX) + " Failed to get telemetryDynamicMessageDescriptor!"); } else { List fields = telemetryDynamicMessageDescriptor.getFields(); if (CollectionUtils.isEmpty(fields)) { - throw new DataValidationException(invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA) + " " + telemetryDynamicMessageDescriptor.getName() + " fields is empty!"); + throw new DataValidationException(DynamicProtoUtils.invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA, EXCEPTION_PREFIX) + " " + telemetryDynamicMessageDescriptor.getName() + " fields is empty!"); } else if (fields.size() == 2) { Descriptors.FieldDescriptor tsFieldDescriptor = telemetryDynamicMessageDescriptor.findFieldByName("ts"); Descriptors.FieldDescriptor valuesFieldDescriptor = telemetryDynamicMessageDescriptor.findFieldByName("values"); if (tsFieldDescriptor != null && valuesFieldDescriptor != null) { if (!Descriptors.FieldDescriptor.Type.MESSAGE.equals(valuesFieldDescriptor.getType())) { - throw new DataValidationException(invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA) + " Field 'values' has invalid data type. Only message type is supported!"); + throw new DataValidationException(DynamicProtoUtils.invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA, EXCEPTION_PREFIX) + " Field 'values' has invalid data type. Only message type is supported!"); } if (!Descriptors.FieldDescriptor.Type.INT64.equals(tsFieldDescriptor.getType())) { - throw new DataValidationException(invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA) + " Field 'ts' has invalid data type. Only int64 type is supported!"); + throw new DataValidationException(DynamicProtoUtils.invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA, EXCEPTION_PREFIX) + " Field 'ts' has invalid data type. Only int64 type is supported!"); } if (!tsFieldDescriptor.hasOptionalKeyword()) { - throw new DataValidationException(invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA) + " Field 'ts' has invalid label. Field 'ts' should have optional keyword!"); + throw new DataValidationException(DynamicProtoUtils.invalidSchemaProvidedMessage(TELEMETRY_PROTO_SCHEMA, EXCEPTION_PREFIX) + " Field 'ts' has invalid label. Field 'ts' should have optional keyword!"); } } } @@ -384,39 +257,39 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator implements TbNode { - private static ObjectMapper mapper = new ObjectMapper(); - private static final String VALUE = "value"; private static final String TS = "ts"; protected C config; + private boolean fetchToData; + private boolean isTellFailureIfAbsent; + private boolean getLatestValueWithTs; @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { this.config = loadGetAttributesNodeConfig(configuration); - mapper.configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false); - mapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true); + this.fetchToData = config.isFetchToData(); + this.getLatestValueWithTs = config.isGetLatestValueWithTs(); + this.isTellFailureIfAbsent = BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true); } protected abstract C loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException; @@ -86,97 +92,110 @@ public abstract class TbAbstractGetAttributesNode> failuresMap = new ConcurrentHashMap<>(); - ListenableFuture> allFutures = Futures.allAsList( - putLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), - putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap, "cs_"), - putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap, "shared_"), - putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap, "ss_") + ListenableFuture>>> allFutures = Futures.allAsList( + getLatestTelemetry(ctx, entityId, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), + getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap), + getAttrAsync(ctx, entityId, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap), + getAttrAsync(ctx, entityId, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap) ); - withCallback(allFutures, i -> { + withCallback(allFutures, futuresList -> { if (!failuresMap.isEmpty()) { throw reportFailures(failuresMap); } - ctx.tellSuccess(msg); + TbMsgMetaData msgMetaData = msg.getMetaData().copy(); + futuresList.stream().filter(Objects::nonNull).forEach(kvEntriesMap -> { + kvEntriesMap.forEach((keyScope, kvEntryList) -> { + String prefix = getPrefix(keyScope); + kvEntryList.forEach(kvEntry -> { + String key = prefix + kvEntry.getKey(); + if (fetchToData) { + JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, key); + } else { + msgMetaData.putValue(key, kvEntry.getValueAsString()); + } + }); + }); + }); + if (fetchToData) { + ctx.tellSuccess(TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msgMetaData, JacksonUtil.toString(msgDataNode))); + } else { + ctx.tellSuccess(TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msgMetaData, msg.getData())); + } }, t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap, String prefix) { + private ListenableFuture>> getAttrAsync(TbContext ctx, EntityId entityId, String scope, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } ListenableFuture> attributeKvEntryListFuture = ctx.getAttributesService().find(ctx.getTenantId(), entityId, scope, keys); return Futures.transform(attributeKvEntryListFuture, attributeKvEntryList -> { - if (!CollectionUtils.isEmpty(attributeKvEntryList)) { - List existingAttributesKvEntry = attributeKvEntryList.stream().filter(attributeKvEntry -> keys.contains(attributeKvEntry.getKey())).collect(Collectors.toList()); - existingAttributesKvEntry.forEach(kvEntry -> msg.getMetaData().putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString())); - if (existingAttributesKvEntry.size() != keys.size() && BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { - getNotExistingKeys(existingAttributesKvEntry, keys).forEach(key -> computeFailuresMap(scope, failuresMap, key)); - } - } else { - if (BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { - keys.forEach(key -> computeFailuresMap(scope, failuresMap, key)); - } + if (isTellFailureIfAbsent && attributeKvEntryList.size() != keys.size()) { + getNotExistingKeys(attributeKvEntryList, keys).forEach(key -> computeFailuresMap(scope, failuresMap, key)); } - return null; + Map> mapAttributeKvEntry = new HashMap<>(); + mapAttributeKvEntry.put(scope, attributeKvEntryList); + return mapAttributeKvEntry; }, MoreExecutors.directExecutor()); } - private ListenableFuture putLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } - ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, keys); - return Futures.transform(latest, l -> { - l.forEach(r -> { - boolean getLatestValueWithTs = BooleanUtils.toBooleanDefaultIfNull(this.config.isGetLatestValueWithTs(), false); - if (BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { - if (r.getValue() == null) { - computeFailuresMap(scope, failuresMap, r.getKey()); - } else if (getLatestValueWithTs) { - putValueWithTs(msg, r); - } else { - msg.getMetaData().putValue(r.getKey(), r.getValueAsString()); + ListenableFuture> latestTelemetryFutures = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, keys); + return Futures.transform(latestTelemetryFutures, tsKvEntries -> { + List listTsKvEntry = new ArrayList<>(); + tsKvEntries.forEach(tsKvEntry -> { + if (tsKvEntry.getValue() == null) { + if (isTellFailureIfAbsent) { + computeFailuresMap(LATEST_TS, failuresMap, tsKvEntry.getKey()); } + } else if (getLatestValueWithTs) { + listTsKvEntry.add(getValueWithTs(tsKvEntry)); } else { - if (r.getValue() != null) { - if (getLatestValueWithTs) { - putValueWithTs(msg, r); - } else { - msg.getMetaData().putValue(r.getKey(), r.getValueAsString()); - } - } + listTsKvEntry.add(new BasicTsKvEntry(tsKvEntry.getTs(), tsKvEntry)); } }); - return null; + Map> mapTsKvEntry = new HashMap<>(); + mapTsKvEntry.put(LATEST_TS, listTsKvEntry); + return mapTsKvEntry; }, MoreExecutors.directExecutor()); } - private void putValueWithTs(TbMsg msg, TsKvEntry r) { - ObjectNode value = mapper.createObjectNode(); - value.put(TS, r.getTs()); - switch (r.getDataType()) { - case STRING: - value.put(VALUE, r.getValueAsString()); - break; - case LONG: - value.put(VALUE, r.getLongValue().get()); - break; - case BOOLEAN: - value.put(VALUE, r.getBooleanValue().get()); + private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) { + ObjectMapper mapper = fetchToData ? JacksonUtil.OBJECT_MAPPER : JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER; + ObjectNode value = JacksonUtil.newObjectNode(mapper); + value.put(TS, tsKvEntry.getTs()); + JacksonUtil.addKvEntry(value, tsKvEntry, VALUE, mapper); + return new BasicTsKvEntry(tsKvEntry.getTs(), new JsonDataEntry(tsKvEntry.getKey(), value.toString())); + } + + private String getPrefix(String scope) { + String prefix = ""; + switch (scope) { + case CLIENT_SCOPE: + prefix = "cs_"; break; - case DOUBLE: - value.put(VALUE, r.getDoubleValue().get()); + case SHARED_SCOPE: + prefix = "shared_"; break; - case JSON: - try { - value.set(VALUE, mapper.readTree(r.getJsonValue().get())); - } catch (IOException e) { - throw new JsonParseException("Can't parse jsonValue: " + r.getJsonValue().get(), e); - } + case SERVER_SCOPE: + prefix = "ss_"; break; } - msg.getMetaData().putValue(r.getKey(), value.toString()); + return prefix; } private List getNotExistingKeys(List existingAttributesKvEntry, List allKeys) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java index 431fe64c09..472b84c803 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java @@ -34,9 +34,9 @@ import org.thingsboard.server.common.msg.TbMsg; @RuleNode(type = ComponentType.ENRICHMENT, name = "originator attributes", configClazz = TbGetAttributesNodeConfiguration.class, - nodeDescription = "Add Message Originator Attributes or Latest Telemetry into Message Metadata", - nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message metadata " + - "with specific prefix: cs/shared/ss. Latest telemetry value added into metadata without prefix. " + + nodeDescription = "Add Message Originator Attributes or Latest Telemetry into Message Data or Metadata", + nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message data/metadata " + + "with specific prefix: cs/shared/ss. Latest telemetry value added into Message data/metadata without prefix. " + "To access those attributes in other nodes this template can be used " + "metadata.cs_temperature or metadata.shared_limit ", uiResources = {"static/rulenode/rulenode-core-config.js"}, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java index 4fc892296f..67766e5b51 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java @@ -35,6 +35,7 @@ public class TbGetAttributesNodeConfiguration implements NodeConfigurationcs/shared/ss. Latest telemetry value added into metadata without prefix. " + + nodeDescription = "Add Originators Related Device Attributes and Latest Telemetry value into Message Data or Metadata", + nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message data/metadata " + + "with specific prefix: cs/shared/ss. Latest telemetry value added into Message data/metadata without prefix. " + "To access those attributes in other nodes this template can be used " + "metadata.cs_temperature or metadata.shared_limit ", uiResources = {"static/rulenode/rulenode-core-config.js"}, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java index 646e2dfef0..049fe39367 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java @@ -36,6 +36,7 @@ public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfigu configuration.setLatestTsKeyNames(Collections.emptyList()); configuration.setTellFailureIfAbsent(true); configuration.setGetLatestValueWithTs(false); + configuration.setFetchToData(false); DeviceRelationsQuery deviceRelationsQuery = new DeviceRelationsQuery(); deviceRelationsQuery.setDirection(EntitySearchDirection.FROM); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java index 77298cba77..045555eae8 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java @@ -15,25 +15,22 @@ */ package org.thingsboard.rule.engine.metadata; -import com.fasterxml.jackson.core.JsonParser; -import com.fasterxml.jackson.core.json.JsonWriteFeature; -import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.ListenableFuture; -import com.google.gson.JsonParseException; import lombok.Data; import lombok.NoArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.StringUtils; import org.apache.commons.lang3.math.NumberUtils; import org.thingsboard.common.util.DonAsynchron; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; @@ -41,7 +38,6 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; -import java.io.IOException; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -75,7 +71,6 @@ public class TbGetTelemetryNode implements TbNode { private TbGetTelemetryNodeConfiguration config; private List tsKeyNames; private int limit; - private ObjectMapper mapper; private String fetchMode; private String orderByFetchAll; private Aggregation aggregation; @@ -91,10 +86,6 @@ public class TbGetTelemetryNode implements TbNode { orderByFetchAll = ASC_ORDER; } aggregation = parseAggregationConfig(config.getAggregation()); - - mapper = new ObjectMapper(); - mapper.configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false); - mapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true); } Aggregation parseAggregationConfig(String aggName) { @@ -146,7 +137,7 @@ public class TbGetTelemetryNode implements TbNode { } private void process(List entries, TbMsg msg, List keys) { - ObjectNode resultNode = mapper.createObjectNode(); + ObjectNode resultNode = JacksonUtil.newObjectNode(JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); if (FETCH_MODE_ALL.equals(fetchMode)) { entries.forEach(entry -> processArray(resultNode, entry)); } else { @@ -169,36 +160,16 @@ public class TbGetTelemetryNode implements TbNode { ArrayNode arrayNode = (ArrayNode) node.get(entry.getKey()); arrayNode.add(buildNode(entry)); } else { - ArrayNode arrayNode = mapper.createArrayNode(); + ArrayNode arrayNode = JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER.createArrayNode(); arrayNode.add(buildNode(entry)); node.set(entry.getKey(), arrayNode); } } private ObjectNode buildNode(TsKvEntry entry) { - ObjectNode obj = mapper.createObjectNode() - .put("ts", entry.getTs()); - switch (entry.getDataType()) { - case STRING: - obj.put("value", entry.getValueAsString()); - break; - case LONG: - obj.put("value", entry.getLongValue().get()); - break; - case BOOLEAN: - obj.put("value", entry.getBooleanValue().get()); - break; - case DOUBLE: - obj.put("value", entry.getDoubleValue().get()); - break; - case JSON: - try { - obj.set("value", mapper.readTree(entry.getJsonValue().get())); - } catch (IOException e) { - throw new JsonParseException("Can't parse jsonValue: " + entry.getJsonValue().get(), e); - } - break; - } + ObjectNode obj = JacksonUtil.newObjectNode(JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); + obj.put("ts", entry.getTs()); + JacksonUtil.addKvEntry(obj, entry, "value", JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); return obj; } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java new file mode 100644 index 0000000000..80554445ed --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java @@ -0,0 +1,338 @@ +/** + * Copyright © 2016-2022 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.rule.engine.metadata; + +import com.datastax.oss.driver.api.core.uuid.Uuids; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.Futures; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; +import org.thingsboard.common.util.AbstractListeningExecutor; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.api.TbContext; +import org.thingsboard.rule.engine.api.TbNodeConfiguration; +import org.thingsboard.rule.engine.api.TbNodeException; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.JsonDataEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; + +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Collectors; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.never; + +@RunWith(MockitoJUnitRunner.class) +public class TbAbstractGetAttributesNodeTest { + + final ObjectMapper mapper = new ObjectMapper(); + + private EntityId originator = new DeviceId(Uuids.timeBased()); + private TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); + + @Mock + private TbContext ctx; + @Mock + private AttributesService attributesService; + @Mock + private TimeseriesService tsService; + private AbstractListeningExecutor dbExecutor; + + private List clientAttributes; + private List serverAttributes; + private List sharedAttributes; + private List tsKeys; + private long ts; + + @Before + public void before() throws TbNodeException { + dbExecutor = new AbstractListeningExecutor() { + @Override + protected int getThreadPollSize() { + return 3; + } + }; + dbExecutor.init(); + + Mockito.reset(ctx); + Mockito.reset(attributesService); + Mockito.reset(tsService); + + Mockito.reset(ctx); + Mockito.reset(attributesService); + Mockito.reset(tsService); + + lenient().when(ctx.getAttributesService()).thenReturn(attributesService); + lenient().when(ctx.getTimeseriesService()).thenReturn(tsService); + lenient().when(ctx.getTenantId()).thenReturn(tenantId); + lenient().when(ctx.getDbCallbackExecutor()).thenReturn(dbExecutor); + + clientAttributes = getAttributeNames("client"); + serverAttributes = getAttributeNames("server"); + sharedAttributes = getAttributeNames("shared"); + tsKeys = List.of("temperature", "humidity", "unknown"); + ts = System.currentTimeMillis(); + + Mockito.when(attributesService.find(tenantId, originator, DataConstants.CLIENT_SCOPE, clientAttributes)) + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(clientAttributes, ts))); + + + Mockito.when(attributesService.find(tenantId, originator, DataConstants.SERVER_SCOPE, serverAttributes)) + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(serverAttributes, ts))); + + + Mockito.when(attributesService.find(tenantId, originator, DataConstants.SHARED_SCOPE, sharedAttributes)) + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(sharedAttributes, ts))); + + Mockito.when(tsService.findLatest(tenantId, originator, tsKeys)) + .thenReturn(Futures.immediateFuture(getListTsKvEntry(tsKeys, ts))); + } + + @After + public void after() { + dbExecutor.destroy(); + } + + @Test + public void fetchToMetadata_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(false, false, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + TbMsgMetaData msgMetaData = resultMsg.getMetaData(); + + //check attributes + checkAttributes(clientAttributes, "cs_", false, msgMetaData, null); + checkAttributes(serverAttributes, "ss_", false, msgMetaData, null); + checkAttributes(sharedAttributes, "shared_", false, msgMetaData, null); + + //check timeseries + checkTs(tsKeys, false, false, msgMetaData, null); + } + + @Test + public void fetchToMetadata_latestWithTs_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(false, true, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + TbMsgMetaData msgMetaData = resultMsg.getMetaData(); + + //check attributes + checkAttributes(clientAttributes, "cs_", false, msgMetaData, null); + checkAttributes(serverAttributes, "ss_", false, msgMetaData, null); + checkAttributes(sharedAttributes, "shared_", false, msgMetaData, null); + + //check timeseries with ts + checkTs(tsKeys, false, true, msgMetaData, null); + } + + @Test + public void fetchToData_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(true, false, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + JsonNode msgData = JacksonUtil.toJsonNode(resultMsg.getData()); + + //check attributes + checkAttributes(clientAttributes, "cs_", true, null, msgData); + checkAttributes(serverAttributes, "ss_", true, null, msgData); + checkAttributes(sharedAttributes, "shared_", true, null, msgData); + + //check timeseries + checkTs(tsKeys, true, false, null, msgData); + } + + @Test + public void fetchToData_latestWithTs_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(true, true, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + JsonNode msgData = JacksonUtil.toJsonNode(resultMsg.getData()); + + //check attributes + checkAttributes(clientAttributes, "cs_", true, null, msgData); + checkAttributes(serverAttributes, "ss_", true, null, msgData); + checkAttributes(sharedAttributes, "shared_", true, null, msgData); + + //check timeseries with ts + checkTs(tsKeys, true, true, null, msgData); + } + + @Test + public void fetchToData_whenOnMsg_then_failure() throws Exception { + TbGetAttributesNode node = initNode(true, true, true); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + ArgumentCaptor newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(Exception.class); + Mockito.verify(ctx, never()).tellSuccess(any()); + Mockito.verify(ctx, Mockito.timeout(5000)).tellFailure(newMsgCaptor.capture(), exceptionCaptor.capture()); + + Assert.assertSame(newMsgCaptor.getValue(), msg); + Assert.assertNotNull(exceptionCaptor.getValue()); + } + + @Test + public void fetchToData_whenOnMsg_then_data_not_object_failure() throws Exception { + TbGetAttributesNode node = initNode(true, true, true); + TbMsg msg = TbMsg.newMsg("TEST", originator, new TbMsgMetaData(), "[]"); + node.onMsg(ctx, msg); + + ArgumentCaptor newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(Exception.class); + Mockito.verify(ctx, never()).tellSuccess(any()); + Mockito.verify(ctx, Mockito.timeout(5000)).tellFailure(newMsgCaptor.capture(), exceptionCaptor.capture()); + + Assert.assertSame(newMsgCaptor.getValue(), msg); + Assert.assertNotNull(exceptionCaptor.getValue()); + } + + private TbMsg checkMsg() { + ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); + Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + + TbMsg resultMsg = msgCaptor.getValue(); + Assert.assertNotNull(resultMsg); + Assert.assertNotNull(resultMsg.getMetaData()); + Assert.assertNotNull(resultMsg.getData()); + return resultMsg; + } + + private void checkAttributes(List attributes, String prefix, boolean fetchToData, TbMsgMetaData msgMetaData, JsonNode msgData) { + attributes.stream() + .filter(attribute -> !attribute.equals("unknown")) + .forEach(attribute -> { + String result; + if (fetchToData) { + result = msgData.get(prefix + attribute).asText(); + } else { + result = msgMetaData.getValue(prefix + attribute); + } + Assert.assertNotNull(result); + Assert.assertEquals(attribute + "_value", result); + }); + } + + private void checkTs(List tsKeys, boolean fetchToData, boolean getLatestValueWithTs, TbMsgMetaData msgMetaData, JsonNode msgData) { + long value = 1L; + for (String key : tsKeys) { + if (key.equals("unknown")) { + continue; + } + String actualValue; + String expectedValue; + if (getLatestValueWithTs) { + expectedValue = "{\"ts\":" + ts + ",\"value\":{\"data\":" + value + "}}"; + } else { + expectedValue = "{\"data\":" + value + "}"; + } + if (fetchToData) { + actualValue = JacksonUtil.toString(msgData.get(key)); + } else { + actualValue = msgMetaData.getValue(key); + } + Assert.assertNotNull(actualValue); + Assert.assertEquals(expectedValue, actualValue); + value++; + } + } + + private TbGetAttributesNode initNode(boolean fetchToData, boolean getLatestValueWithTs, boolean isTellFailureIfAbsent) throws TbNodeException { + TbGetAttributesNodeConfiguration config = new TbGetAttributesNodeConfiguration(); + config.setClientAttributeNames(List.of("client_attr_1", "client_attr_2", "${client_attr_metadata}", "unknown")); + config.setServerAttributeNames(List.of("server_attr_1", "server_attr_2", "${server_attr_metadata}", "unknown")); + config.setSharedAttributeNames(List.of("shared_attr_1", "shared_attr_2", "$[shared_attr_data]", "unknown")); + config.setLatestTsKeyNames(List.of("temperature", "humidity", "unknown")); + config.setFetchToData(fetchToData); + config.setGetLatestValueWithTs(getLatestValueWithTs); + config.setTellFailureIfAbsent(isTellFailureIfAbsent); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + TbGetAttributesNode node = new TbGetAttributesNode(); + node.init(ctx, nodeConfiguration); + return node; + } + + private TbMsg getTbMsg(EntityId entityId) { + ObjectNode msgData = JacksonUtil.newObjectNode(); + msgData.put("shared_attr_data", "shared_attr_3"); + + TbMsgMetaData msgMetaData = new TbMsgMetaData(); + msgMetaData.putValue("client_attr_metadata", "client_attr_3"); + msgMetaData.putValue("server_attr_metadata", "server_attr_3"); + + return TbMsg.newMsg("TEST", entityId, msgMetaData, msgData.toString()); + } + + private List getAttributeNames(String prefix) { + return List.of(prefix + "_attr_1", prefix + "_attr_2", prefix + "_attr_3", "unknown"); + } + + private List getListAttributeKvEntry(List attributes, long ts) { + return attributes.stream() + .filter(attribute -> !attribute.equals("unknown")) + .map(attribute -> toAttributeKvEntry(ts, attribute)) + .collect(Collectors.toList()); + } + + private BaseAttributeKvEntry toAttributeKvEntry(long ts, String attribute) { + return new BaseAttributeKvEntry(ts, new StringDataEntry(attribute, attribute + "_value")); + } + + private List getListTsKvEntry(List keys, long ts) { + long value = 1L; + List kvEntries = new ArrayList<>(); + for (String key : keys) { + if (key.equals("unknown")) { + continue; + } + String dataValue = "{\"data\":" + value + "}"; + kvEntries.add(new BasicTsKvEntry(ts, new JsonDataEntry(key, dataValue))); + value++; + } + return kvEntries; + } + +} diff --git a/ui-ngx/src/app/core/services/utils.service.ts b/ui-ngx/src/app/core/services/utils.service.ts index b901b87328..66eefa9e13 100644 --- a/ui-ngx/src/app/core/services/utils.service.ts +++ b/ui-ngx/src/app/core/services/utils.service.ts @@ -43,6 +43,15 @@ import { WidgetInfo } from '@home/models/widget-component.models'; import jsonSchemaDefaults from 'json-schema-defaults'; import materialIconsCodepoints from '!raw-loader!./material-icons-codepoints.raw'; import { Observable, of, ReplaySubject } from 'rxjs'; +import { publishReplay, refCount } from 'rxjs/operators'; +import { WidgetContext } from '@app/modules/home/models/widget-component.models'; +import { + AttributeData, + LatestTelemetry, + TelemetrySubscriber, + TelemetryType +} from '@shared/models/telemetry/telemetry.models'; +import { EntityId } from '@shared/models/id/entity-id'; const i18nRegExp = new RegExp(`{${i18nPrefix}:[^{}]+}`, 'g'); @@ -479,4 +488,27 @@ export class UtilsService { return defaultValue; } } + + private getEntityIdFromDatasource(dataSource: Datasource): EntityId { + return {id: dataSource.entityId, entityType: dataSource.entityType}; + } + + public subscribeToEntityTelemetry(ctx: WidgetContext, + entityId?: EntityId, + type: TelemetryType = LatestTelemetry.LATEST_TELEMETRY, + keys: string[] = null): Observable> { + if (!entityId && ctx.datasources.length > 0) { + entityId = this.getEntityIdFromDatasource(ctx.datasources[0]); + } + const subscription = TelemetrySubscriber.createEntityAttributesSubscription(ctx.telemetryWsService, entityId, type, ctx.ngZone, keys); + if (!ctx.telemetrySubscribers) { + ctx.telemetrySubscribers = []; + } + ctx.telemetrySubscribers.push(subscription); + subscription.subscribe(); + return subscription.attributeData$().pipe( + publishReplay(1), + refCount() + ); + } } diff --git a/ui-ngx/src/app/modules/home/components/widget/dynamic-widget.component.ts b/ui-ngx/src/app/modules/home/components/widget/dynamic-widget.component.ts index 59b1959269..040e27b2ff 100644 --- a/ui-ngx/src/app/modules/home/components/widget/dynamic-widget.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/dynamic-widget.component.ts @@ -40,6 +40,7 @@ import { AuthService } from '@core/auth/auth.service'; import { DialogService } from '@core/services/dialog.service'; import { CustomDialogService } from '@home/components/widget/dialog/custom-dialog.service'; import { ResourceService } from '@core/http/resource.service'; +import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; import { DatePipe } from '@angular/common'; import { TranslateService } from '@ngx-translate/core'; import { DomSanitizer } from '@angular/platform-browser'; @@ -80,6 +81,7 @@ export class DynamicWidgetComponent extends PageComponent implements IDynamicWid this.ctx.dialogs = $injector.get(DialogService); this.ctx.customDialog = $injector.get(CustomDialogService); this.ctx.resourceService = $injector.get(ResourceService); + this.ctx.telemetryWsService = $injector.get(TelemetryWebsocketService); this.ctx.date = $injector.get(DatePipe); this.ctx.translate = $injector.get(TranslateService); this.ctx.http = $injector.get(HttpClient); @@ -100,7 +102,9 @@ export class DynamicWidgetComponent extends PageComponent implements IDynamicWid } ngOnDestroy(): void { - + if (this.ctx.telemetrySubscribers) { + this.ctx.telemetrySubscribers.forEach(item => item.unsubscribe()); + } } clearRpcError() { diff --git a/ui-ngx/src/app/modules/home/models/services.map.ts b/ui-ngx/src/app/modules/home/models/services.map.ts index 93ec69da61..7514e972b5 100644 --- a/ui-ngx/src/app/modules/home/models/services.map.ts +++ b/ui-ngx/src/app/modules/home/models/services.map.ts @@ -39,6 +39,7 @@ import { OtaPackageService } from '@core/http/ota-package.service'; import { AuthService } from '@core/auth/auth.service'; import { ResourceService } from '@core/http/resource.service'; import { TwoFactorAuthenticationService } from '@core/http/two-factor-authentication.service'; +import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; export const ServicesMap = new Map>( [ @@ -65,6 +66,7 @@ export const ServicesMap = new Map>( ['otaPackageService', OtaPackageService], ['authService', AuthService], ['resourceService', ResourceService], - ['twoFactorAuthenticationService', TwoFactorAuthenticationService] + ['twoFactorAuthenticationService', TwoFactorAuthenticationService], + ['telemetryWsService', TelemetryWebsocketService] ] ); diff --git a/ui-ngx/src/app/modules/home/models/widget-component.models.ts b/ui-ngx/src/app/modules/home/models/widget-component.models.ts index afd50d1c82..813b3b6927 100644 --- a/ui-ngx/src/app/modules/home/models/widget-component.models.ts +++ b/ui-ngx/src/app/modules/home/models/widget-component.models.ts @@ -75,6 +75,7 @@ import { DialogService } from '@core/services/dialog.service'; import { CustomDialogService } from '@home/components/widget/dialog/custom-dialog.service'; import { AuthService } from '@core/auth/auth.service'; import { ResourceService } from '@core/http/resource.service'; +import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; import { DatePipe } from '@angular/common'; import { TranslateService } from '@ngx-translate/core'; import { PageLink, TimePageLink } from '@shared/models/page/page-link'; @@ -87,6 +88,7 @@ import * as RxJSOperators from 'rxjs/operators'; import { TbPopoverComponent } from '@shared/components/popover.component'; import { EntityId } from '@shared/models/id/entity-id'; import { AlarmQuery, AlarmSearchStatus, AlarmStatus} from '@app/shared/models/alarm.models'; +import { TelemetrySubscriber } from '@app/shared/public-api'; export interface IWidgetAction { name: string; @@ -177,6 +179,8 @@ export class WidgetContext { dialogs: DialogService; customDialog: CustomDialogService; resourceService: ResourceService; + telemetryWsService: TelemetryWebsocketService; + telemetrySubscribers?: TelemetrySubscriber[]; date: DatePipe; translate: TranslateService; http: HttpClient;