From c9966ef3ef4ce1f5be77109103de2ff9f8416d27 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 2 Jun 2023 16:27:06 +0300 Subject: [PATCH] Extracted proto files to a separate module --- .../device/DeviceActorMessageProcessor.java | 27 +-- .../queue/DefaultTbCoreConsumerService.java | 5 +- .../subscription/TbSubscriptionUtils.java | 65 +------- common/cluster-api/pom.xml | 8 +- common/pom.xml | 1 + common/proto/pom.xml | 109 ++++++++++++ .../transport/adaptor/AdaptorException.java | 0 .../transport/adaptor/JsonConverter.java | 0 .../adaptor/JsonConverterConfig.java | 0 .../transport/adaptor/ProtoConverter.java | 0 .../common/transport/util/KvProtoUtil.java | 156 ++++++++++++++++++ .../src/main/proto/jsinvoke.proto | 0 .../src/main/proto/queue.proto | 0 .../src/main/proto/transport.proto | 0 common/queue/pom.xml | 4 + common/transport/transport-api/pom.xml | 9 - pom.xml | 7 + 17 files changed, 292 insertions(+), 99 deletions(-) create mode 100644 common/proto/pom.xml rename common/{transport/transport-api => proto}/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java (100%) rename common/{transport/transport-api => proto}/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java (100%) rename common/{transport/transport-api => proto}/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverterConfig.java (100%) rename common/{transport/transport-api => proto}/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java (100%) create mode 100644 common/proto/src/main/java/org/thingsboard/server/common/transport/util/KvProtoUtil.java rename common/{cluster-api => proto}/src/main/proto/jsinvoke.proto (100%) rename common/{cluster-api => proto}/src/main/proto/queue.proto (100%) rename common/{transport/transport-api => proto}/src/main/proto/transport.proto (100%) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 934f309480..1831c6ab9c 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -63,6 +63,7 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg; +import org.thingsboard.server.common.transport.util.KvProtoUtil; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; @@ -455,7 +456,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso GetAttributeResponseMsg responseMsg = GetAttributeResponseMsg.newBuilder() .setRequestId(requestId) .setSharedStateMsg(true) - .addAllSharedAttributeList(toTsKvProtos(result)) + .addAllSharedAttributeList(KvProtoUtil.attrToTsKvProtos(result)) .setIsMultipleAttributesRequest(request.getSharedAttributeNamesCount() > 1) .build(); sendToTransport(responseMsg, sessionInfo); @@ -476,8 +477,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso public void onSuccess(@Nullable List> result) { GetAttributeResponseMsg responseMsg = GetAttributeResponseMsg.newBuilder() .setRequestId(requestId) - .addAllClientAttributeList(toTsKvProtos(result.get(0))) - .addAllSharedAttributeList(toTsKvProtos(result.get(1))) + .addAllClientAttributeList(KvProtoUtil.attrToTsKvProtos(result.get(0))) + .addAllSharedAttributeList(KvProtoUtil.attrToTsKvProtos(result.get(1))) .setIsMultipleAttributesRequest( request.getSharedAttributeNamesCount() + request.getClientAttributeNamesCount() > 1) .build(); @@ -547,7 +548,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) { List attributes = new ArrayList<>(msg.getValues()); if (attributes.size() > 0) { - List sharedUpdated = msg.getValues().stream().map(this::toTsKvProto) + List sharedUpdated = msg.getValues().stream().map(t -> KvProtoUtil.toTsKvProto(t.getLastUpdateTs(), t)) .collect(Collectors.toList()); if (!sharedUpdated.isEmpty()) { notification.addAllSharedUpdated(sharedUpdated); @@ -834,24 +835,6 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso }, systemContext.getDbCallbackExecutor()); } - private List toTsKvProtos(@Nullable List result) { - List clientAttributes; - if (result == null || result.isEmpty()) { - clientAttributes = Collections.emptyList(); - } else { - clientAttributes = new ArrayList<>(result.size()); - for (AttributeKvEntry attrEntry : result) { - clientAttributes.add(toTsKvProto(attrEntry)); - } - } - return clientAttributes; - } - - private TsKvProto toTsKvProto(AttributeKvEntry attrEntry) { - return TsKvProto.newBuilder().setTs(attrEntry.getLastUpdateTs()) - .setKv(toKeyValueProto(attrEntry)).build(); - } - private KeyValueProto toKeyValueProto(KvEntry kvEntry) { KeyValueProto.Builder builder = KeyValueProto.newBuilder(); builder.setKey(kvEntry.getKey()); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 2bd0043046..e11b6cdd5b 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -40,6 +40,7 @@ import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.stats.StatsFactory; +import org.thingsboard.server.common.transport.util.KvProtoUtil; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto; @@ -522,13 +523,13 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService builder.addData(toKeyValueProto(v.getTs(), v).build())); + ts.forEach(v -> builder.addData(KvProtoUtil.toKeyValueProto(v.getTs(), v).build())); SubscriptionMgrMsgProto.Builder msgBuilder = SubscriptionMgrMsgProto.newBuilder(); msgBuilder.setTsUpdate(builder); return ToCoreMsg.newBuilder().setToSubscriptionMgrMsg(msgBuilder.build()).build(); @@ -269,7 +270,7 @@ public class TbSubscriptionUtils { builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); builder.setScope(scope); - attributes.forEach(v -> builder.addData(toKeyValueProto(v.getLastUpdateTs(), v).build())); + attributes.forEach(v -> builder.addData(KvProtoUtil.toKeyValueProto(v.getLastUpdateTs(), v).build())); SubscriptionMgrMsgProto.Builder msgBuilder = SubscriptionMgrMsgProto.newBuilder(); msgBuilder.setAttrUpdate(builder); @@ -292,70 +293,10 @@ public class TbSubscriptionUtils { return ToCoreMsg.newBuilder().setToSubscriptionMgrMsg(msgBuilder.build()).build(); } - - private static TsKvProto.Builder toKeyValueProto(long ts, KvEntry attr) { - KeyValueProto.Builder dataBuilder = KeyValueProto.newBuilder(); - dataBuilder.setKey(attr.getKey()); - dataBuilder.setType(KeyValueType.forNumber(attr.getDataType().ordinal())); - switch (attr.getDataType()) { - case BOOLEAN: - attr.getBooleanValue().ifPresent(dataBuilder::setBoolV); - break; - case LONG: - attr.getLongValue().ifPresent(dataBuilder::setLongV); - break; - case DOUBLE: - attr.getDoubleValue().ifPresent(dataBuilder::setDoubleV); - break; - case JSON: - attr.getJsonValue().ifPresent(dataBuilder::setJsonV); - break; - case STRING: - attr.getStrValue().ifPresent(dataBuilder::setStringV); - break; - } - return TsKvProto.newBuilder().setTs(ts).setKv(dataBuilder); - } - public static EntityId toEntityId(String entityType, long entityIdMSB, long entityIdLSB) { return EntityIdFactory.getByTypeAndUuid(entityType, new UUID(entityIdMSB, entityIdLSB)); } - public static List toTsKvEntityList(List dataList) { - List result = new ArrayList<>(dataList.size()); - dataList.forEach(proto -> result.add(new BasicTsKvEntry(proto.getTs(), getKvEntry(proto.getKv())))); - return result; - } - - public static List toAttributeKvList(List dataList) { - List result = new ArrayList<>(dataList.size()); - dataList.forEach(proto -> result.add(new BaseAttributeKvEntry(getKvEntry(proto.getKv()), proto.getTs()))); - return result; - } - - private static KvEntry getKvEntry(KeyValueProto proto) { - KvEntry entry = null; - DataType type = DataType.values()[proto.getType().getNumber()]; - switch (type) { - case BOOLEAN: - entry = new BooleanDataEntry(proto.getKey(), proto.getBoolV()); - break; - case LONG: - entry = new LongDataEntry(proto.getKey(), proto.getLongV()); - break; - case DOUBLE: - entry = new DoubleDataEntry(proto.getKey(), proto.getDoubleV()); - break; - case STRING: - entry = new StringDataEntry(proto.getKey(), proto.getStringV()); - break; - case JSON: - entry = new JsonDataEntry(proto.getKey(), proto.getJsonV()); - break; - } - return entry; - } - public static ToCoreMsg toAlarmUpdateProto(TenantId tenantId, EntityId entityId, AlarmInfo alarm) { TbAlarmUpdateProto.Builder builder = TbAlarmUpdateProto.newBuilder(); builder.setEntityType(entityId.getEntityType().name()); diff --git a/common/cluster-api/pom.xml b/common/cluster-api/pom.xml index c9fa6d0894..93cd7c665a 100644 --- a/common/cluster-api/pom.xml +++ b/common/cluster-api/pom.xml @@ -40,6 +40,10 @@ org.thingsboard.common data + + org.thingsboard.common + proto + org.thingsboard.common message @@ -119,10 +123,6 @@ - - org.xolstice.maven.plugins - protobuf-maven-plugin - org.apache.maven.plugins maven-source-plugin diff --git a/common/pom.xml b/common/pom.xml index 0edd153a3c..0c80189acf 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -35,6 +35,7 @@ data + proto util message actor diff --git a/common/proto/pom.xml b/common/proto/pom.xml new file mode 100644 index 0000000000..ab2cdf91c9 --- /dev/null +++ b/common/proto/pom.xml @@ -0,0 +1,109 @@ + + + 4.0.0 + + org.thingsboard + 3.6.0-SNAPSHOT + common + + org.thingsboard.common + proto + jar + + Thingsboard Server Common Protobuf and gRPC structures + https://thingsboard.io + + + UTF-8 + ${basedir}/../.. + + + + + org.thingsboard.common + data + + + org.thingsboard.common + message + + + com.google.protobuf + protobuf-java + + + com.google.protobuf + protobuf-java-util + + + io.grpc + grpc-netty-shaded + provided + + + io.grpc + grpc-protobuf + provided + + + io.grpc + grpc-stub + provided + + + org.springframework.boot + spring-boot-starter-web + provided + + + org.springframework.boot + spring-boot-starter-test + test + + + org.junit.vintage + junit-vintage-engine + test + + + org.awaitility + awaitility + test + + + + + + + org.xolstice.maven.plugins + protobuf-maven-plugin + + + + + + + thingsboard-repo-deploy + ThingsBoard Repo Deployment + https://repo.thingsboard.io/artifactory/libs-release-public + + + + diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java b/common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java similarity index 100% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java rename to common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java b/common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java similarity index 100% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java rename to common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverterConfig.java b/common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverterConfig.java similarity index 100% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverterConfig.java rename to common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverterConfig.java diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java b/common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java similarity index 100% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java rename to common/proto/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java diff --git a/common/proto/src/main/java/org/thingsboard/server/common/transport/util/KvProtoUtil.java b/common/proto/src/main/java/org/thingsboard/server/common/transport/util/KvProtoUtil.java new file mode 100644 index 0000000000..172d014460 --- /dev/null +++ b/common/proto/src/main/java/org/thingsboard/server/common/transport/util/KvProtoUtil.java @@ -0,0 +1,156 @@ +/** + * Copyright © 2016-2023 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.transport.util; + +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.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DataType; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.JsonDataEntry; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.gen.transport.TransportProtos; + +import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +public class KvProtoUtil { + + public static List attrToTsKvProtos(List result) { + List clientAttributes; + if (result == null || result.isEmpty()) { + clientAttributes = Collections.emptyList(); + } else { + clientAttributes = new ArrayList<>(result.size()); + for (AttributeKvEntry attrEntry : result) { + clientAttributes.add(toTsKvProto(attrEntry.getLastUpdateTs(), attrEntry)); + } + } + return clientAttributes; + } + + + public static List tsToTsKvProtos(List result) { + List ts; + if (result == null || result.isEmpty()) { + ts = Collections.emptyList(); + } else { + ts = new ArrayList<>(result.size()); + for (TsKvEntry attrEntry : result) { + ts.add(toTsKvProto(attrEntry.getTs(), attrEntry)); + } + } + return ts; + } + + public static TransportProtos.TsKvProto toTsKvProto(long ts, KvEntry kvEntry) { + return TransportProtos.TsKvProto.newBuilder().setTs(ts) + .setKv(KvProtoUtil.toKeyValueProto(kvEntry)).build(); + } + + public static TransportProtos.KeyValueProto toKeyValueProto(KvEntry kvEntry) { + TransportProtos.KeyValueProto.Builder builder = TransportProtos.KeyValueProto.newBuilder(); + builder.setKey(kvEntry.getKey()); + switch (kvEntry.getDataType()) { + case BOOLEAN: + builder.setType(TransportProtos.KeyValueType.BOOLEAN_V); + builder.setBoolV(kvEntry.getBooleanValue().get()); + break; + case DOUBLE: + builder.setType(TransportProtos.KeyValueType.DOUBLE_V); + builder.setDoubleV(kvEntry.getDoubleValue().get()); + break; + case LONG: + builder.setType(TransportProtos.KeyValueType.LONG_V); + builder.setLongV(kvEntry.getLongValue().get()); + break; + case STRING: + builder.setType(TransportProtos.KeyValueType.STRING_V); + builder.setStringV(kvEntry.getStrValue().get()); + break; + case JSON: + builder.setType(TransportProtos.KeyValueType.JSON_V); + builder.setJsonV(kvEntry.getJsonValue().get()); + break; + } + return builder.build(); + } + + public static TransportProtos.TsKvProto.Builder toKeyValueProto(long ts, KvEntry attr) { + TransportProtos.KeyValueProto.Builder dataBuilder = TransportProtos.KeyValueProto.newBuilder(); + dataBuilder.setKey(attr.getKey()); + dataBuilder.setType(TransportProtos.KeyValueType.forNumber(attr.getDataType().ordinal())); + switch (attr.getDataType()) { + case BOOLEAN: + attr.getBooleanValue().ifPresent(dataBuilder::setBoolV); + break; + case LONG: + attr.getLongValue().ifPresent(dataBuilder::setLongV); + break; + case DOUBLE: + attr.getDoubleValue().ifPresent(dataBuilder::setDoubleV); + break; + case JSON: + attr.getJsonValue().ifPresent(dataBuilder::setJsonV); + break; + case STRING: + attr.getStrValue().ifPresent(dataBuilder::setStringV); + break; + } + return TransportProtos.TsKvProto.newBuilder().setTs(ts).setKv(dataBuilder); + } + + public static List toTsKvEntityList(List dataList) { + List result = new ArrayList<>(dataList.size()); + dataList.forEach(proto -> result.add(new BasicTsKvEntry(proto.getTs(), getKvEntry(proto.getKv())))); + return result; + } + + public static List toAttributeKvList(List dataList) { + List result = new ArrayList<>(dataList.size()); + dataList.forEach(proto -> result.add(new BaseAttributeKvEntry(getKvEntry(proto.getKv()), proto.getTs()))); + return result; + } + + private static KvEntry getKvEntry(TransportProtos.KeyValueProto proto) { + KvEntry entry = null; + DataType type = DataType.values()[proto.getType().getNumber()]; + switch (type) { + case BOOLEAN: + entry = new BooleanDataEntry(proto.getKey(), proto.getBoolV()); + break; + case LONG: + entry = new LongDataEntry(proto.getKey(), proto.getLongV()); + break; + case DOUBLE: + entry = new DoubleDataEntry(proto.getKey(), proto.getDoubleV()); + break; + case STRING: + entry = new StringDataEntry(proto.getKey(), proto.getStringV()); + break; + case JSON: + entry = new JsonDataEntry(proto.getKey(), proto.getJsonV()); + break; + } + return entry; + } +} diff --git a/common/cluster-api/src/main/proto/jsinvoke.proto b/common/proto/src/main/proto/jsinvoke.proto similarity index 100% rename from common/cluster-api/src/main/proto/jsinvoke.proto rename to common/proto/src/main/proto/jsinvoke.proto diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto similarity index 100% rename from common/cluster-api/src/main/proto/queue.proto rename to common/proto/src/main/proto/queue.proto diff --git a/common/transport/transport-api/src/main/proto/transport.proto b/common/proto/src/main/proto/transport.proto similarity index 100% rename from common/transport/transport-api/src/main/proto/transport.proto rename to common/proto/src/main/proto/transport.proto diff --git a/common/queue/pom.xml b/common/queue/pom.xml index edd4574bdd..59e96d210d 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -36,6 +36,10 @@ + + org.thingsboard.common + proto + org.thingsboard.common data diff --git a/common/transport/transport-api/pom.xml b/common/transport/transport-api/pom.xml index 2e2174e0a9..b45cd28e52 100644 --- a/common/transport/transport-api/pom.xml +++ b/common/transport/transport-api/pom.xml @@ -135,13 +135,4 @@ - - - - org.xolstice.maven.plugins - protobuf-maven-plugin - - - - diff --git a/pom.xml b/pom.xml index 9b11160fab..c778ebcaab 100755 --- a/pom.xml +++ b/pom.xml @@ -28,6 +28,8 @@ 2016 + 17 + 17 ${basedir} true none @@ -892,6 +894,11 @@ version-control ${project.version} + + org.thingsboard.common + proto + ${project.version} + org.thingsboard.common cache