diff --git a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java index e7243ed3bc..5f30bc291a 100644 --- a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java @@ -55,9 +55,9 @@ import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.msg.edqs.EdqsService; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.edqs.processor.EdqsConverter; +import org.thingsboard.server.edqs.util.EdqsConverter; import org.thingsboard.server.edqs.processor.EdqsProducer; -import org.thingsboard.server.edqs.util.EdqsPartitionService; +import org.thingsboard.server.edqs.state.EdqsPartitionService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg; import org.thingsboard.server.gen.transport.TransportProtos.EdqsRequestMsg; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/AttributeKv.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/AttributeKv.java index 9f65f72be4..c0e3abe9e3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/AttributeKv.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/AttributeKv.java @@ -20,6 +20,7 @@ import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; @@ -66,4 +67,9 @@ public class AttributeKv implements EdqsObject { return version; } + @Override + public ObjectType type() { + return ObjectType.ATTRIBUTE_KV; + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsObject.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsObject.java index 9a2836149a..d1e463443c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsObject.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsObject.java @@ -16,6 +16,7 @@ package org.thingsboard.server.common.data.edqs; import com.fasterxml.jackson.annotation.JsonIgnore; +import org.thingsboard.server.common.data.ObjectType; public interface EdqsObject { @@ -25,4 +26,7 @@ public interface EdqsObject { @JsonIgnore Long version(); + @JsonIgnore + ObjectType type(); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/Entity.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/Entity.java index decfbbd116..c08047c005 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/Entity.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/Entity.java @@ -20,6 +20,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; import lombok.Data; import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.edqs.fields.EntityFields; import org.thingsboard.server.common.data.edqs.fields.EntityIdFields; @@ -59,4 +60,9 @@ public class Entity implements EdqsObject { return fields.getVersion(); } + @Override + public ObjectType type() { + return ObjectType.fromEntityType(type); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/LatestTsKv.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/LatestTsKv.java index 19a767a320..b6a466df12 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/LatestTsKv.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/LatestTsKv.java @@ -19,6 +19,7 @@ import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @@ -61,4 +62,9 @@ public class LatestTsKv implements EdqsObject { return version; } + @Override + public ObjectType type() { + return ObjectType.LATEST_TS_KV; + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java b/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java index 927d592cd7..c04c490a26 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java @@ -25,6 +25,7 @@ import lombok.ToString; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.BaseDataWithAdditionalInfo; import org.thingsboard.server.common.data.HasVersion; +import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.edqs.EdqsObject; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.validation.Length; @@ -127,4 +128,9 @@ public class EntityRelation implements HasVersion, Serializable, EdqsObject { return version; } + @Override + public ObjectType type() { + return ObjectType.RELATION; + } + } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index a90a029bc4..7e9d60d09a 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java @@ -49,7 +49,8 @@ import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.edqs.repo.EdqRepository; import org.thingsboard.server.edqs.state.EdqsStateService; -import org.thingsboard.server.edqs.util.EdqsPartitionService; +import org.thingsboard.server.edqs.util.EdqsConverter; +import org.thingsboard.server.edqs.state.EdqsPartitionService; import org.thingsboard.server.edqs.util.VersionsStore; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg; diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProducer.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProducer.java index 8973f3b5b8..bfd0d9df59 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProducer.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProducer.java @@ -21,7 +21,7 @@ import org.apache.kafka.common.errors.RecordTooLargeException; import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.edqs.util.EdqsPartitionService; +import org.thingsboard.server.edqs.state.EdqsPartitionService; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueMsgMetadata; diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsPartitionService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsPartitionService.java similarity index 97% rename from common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsPartitionService.java rename to common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsPartitionService.java index 45de60687c..3d7757f17b 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsPartitionService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsPartitionService.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.edqs.util; +package org.thingsboard.server.edqs.state; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index cfa831f720..6ab1fb9e1f 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java @@ -30,7 +30,6 @@ import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.edqs.processor.EdqsProcessor; import org.thingsboard.server.edqs.processor.EdqsProducer; -import org.thingsboard.server.edqs.util.EdqsPartitionService; import org.thingsboard.server.edqs.util.VersionsStore; import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsConverter.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java similarity index 99% rename from common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsConverter.java rename to common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java index f5e405e516..b7419451b9 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsConverter.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.edqs.processor; +package org.thingsboard.server.edqs.util; import com.google.protobuf.ByteString; import lombok.RequiredArgsConstructor; diff --git a/edqs/src/test/java/org/thingsboard/server/edqs/repo/AbstractEDQTest.java b/edqs/src/test/java/org/thingsboard/server/edqs/repo/AbstractEDQTest.java index 8fb680df1c..1f22f9d5ae 100644 --- a/edqs/src/test/java/org/thingsboard/server/edqs/repo/AbstractEDQTest.java +++ b/edqs/src/test/java/org/thingsboard/server/edqs/repo/AbstractEDQTest.java @@ -59,7 +59,7 @@ import org.thingsboard.server.common.data.query.KeyFilter; import org.thingsboard.server.common.data.query.StringFilterPredicate; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; -import org.thingsboard.server.edqs.processor.EdqsConverter; +import org.thingsboard.server.edqs.util.EdqsConverter; import java.util.Collections; import java.util.List; @@ -67,7 +67,7 @@ import java.util.UUID; @RunWith(SpringRunner.class) @Configuration -@ComponentScan({"org.thingsboard.server.edqs.repo"}) +@ComponentScan({"org.thingsboard.server.edqs.repo", "org.thingsboard.server.edqs.util"}) @EntityScan("org.thingsboard.server.edqs") @TestPropertySource(locations = {"classpath:edq-test.properties"}) @TestExecutionListeners({ @@ -77,6 +77,8 @@ public abstract class AbstractEDQTest { @Autowired protected InMemoryEdqRepository repository; + @Autowired + protected EdqsConverter edqsConverter; protected final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); protected final CustomerId customerId = new CustomerId(UUID.randomUUID()); @@ -218,7 +220,7 @@ public abstract class AbstractEDQTest { } protected void createRelation(EntityType fromType, UUID fromId, EntityType toType, UUID toId, RelationTypeGroup group, String type) { - repository.get(tenantId).addOrUpdate(new EntityRelation(EntityIdFactory.getByTypeAndUuid(fromType, fromId), EntityIdFactory.getByTypeAndUuid(toType, toId), type, group)); + addOrUpdate(new EntityRelation(EntityIdFactory.getByTypeAndUuid(fromType, fromId), EntityIdFactory.getByTypeAndUuid(toType, toId), type, group)); } @@ -243,6 +245,8 @@ public abstract class AbstractEDQTest { } protected void addOrUpdate(EdqsObject edqsObject) { + byte[] serialized = edqsConverter.serialize(edqsObject.type(), edqsObject); + edqsObject = edqsConverter.deserialize(edqsObject.type(), serialized); repository.get(tenantId).addOrUpdate(edqsObject); }