From cc1887fc8dac658535b2c6c2858f210b7cc3e936 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 26 Jun 2020 11:46:12 +0300 Subject: [PATCH 1/4] Added WS notification on deleted attribute --- .../controller/TelemetryController.java | 7 +++-- .../queue/DefaultTbCoreConsumerService.java | 8 ++++++ .../DefaultSubscriptionManagerService.java | 27 +++++++++++++++++++ .../SubscriptionManagerService.java | 1 + .../subscription/TbEntityDataSubCtx.java | 17 +++++++----- .../subscription/TbSubscriptionUtils.java | 17 ++++++++++++ .../DefaultTelemetrySubscriptionService.java | 21 +++++++++++++++ application/src/main/resources/logback.xml | 2 ++ common/queue/src/main/proto/queue.proto | 11 ++++++++ .../query/DefaultEntityQueryRepository.java | 2 +- .../server/dao/SqlDaoServiceTestSuite.java | 2 +- .../dao/service/BaseEntityServiceTest.java | 13 ++++++--- .../api/RuleEngineTelemetryService.java | 3 +++ 13 files changed, 116 insertions(+), 15 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java index bb866e3e9b..f61565d320 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java @@ -355,10 +355,9 @@ public class TelemetryController extends BaseController { DataConstants.SHARED_SCOPE.equals(scope) || DataConstants.CLIENT_SCOPE.equals(scope)) { return accessValidator.validateEntityAndCallback(getCurrentUser(), Operation.WRITE_ATTRIBUTES, entityIdSrc, (result, tenantId, entityId) -> { - ListenableFuture> future = attributesService.removeAll(user.getTenantId(), entityId, scope, keys); - Futures.addCallback(future, new FutureCallback>() { + tsSubService.deleteAndNotify(tenantId, entityId, scope, keys, new FutureCallback() { @Override - public void onSuccess(@Nullable List tmp) { + public void onSuccess(@Nullable Void tmp) { logAttributesDeleted(user, entityId, scope, keys, null); if (entityIdSrc.getEntityType().equals(EntityType.DEVICE)) { DeviceId deviceId = new DeviceId(entityId.getId()); @@ -375,7 +374,7 @@ public class TelemetryController extends BaseController { logAttributesDeleted(user, entityId, scope, keys, t); result.setResult(new ResponseEntity<>(HttpStatus.INTERNAL_SERVER_ERROR)); } - }, executor); + }); }); } else { return getImmediateDeferredResult("Invalid attribute scope: " + scope, HttpStatus.BAD_REQUEST); 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 ff88f012ee..1d1ba3a3d5 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 @@ -31,6 +31,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.FromDeviceRPCRespons import org.thingsboard.server.gen.transport.TransportProtos.LocalSubscriptionServiceMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionMgrMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeUpdateProto; +import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeDeleteProto; import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionCloseProto; import org.thingsboard.server.gen.transport.TransportProtos.TbTimeSeriesUpdateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -54,6 +55,7 @@ import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWra import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; +import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -259,6 +261,12 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService keys, TbCallback callback) { + onLocalSubUpdate(entityId, + s -> { + if (TbSubscriptionType.ATTRIBUTES.equals(s.getType())) { + return (TbAttributeSubscription) s; + } else { + return null; + } + }, + s -> (TbAttributeSubscriptionScope.ANY_SCOPE.equals(s.getScope()) || scope.equals(s.getScope().name())), + s -> { + List subscriptionUpdate = null; + for (String key : keys) { + if (s.isAllKeys() || s.getKeyStates().containsKey(key)) { + if (subscriptionUpdate == null) { + subscriptionUpdate = new ArrayList<>(); + } + subscriptionUpdate.add(new BasicTsKvEntry(0, new StringDataEntry(key, null))); + } + } + return subscriptionUpdate; + }); + callback.onSuccess(); + } + private void onLocalSubUpdate(EntityId entityId, Function castFunction, Predicate filterFunction, diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java index e20b707348..af868374aa 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java @@ -35,4 +35,5 @@ public interface SubscriptionManagerService extends ApplicationListener attributes, TbCallback callback); + void onAttributesDelete(TenantId tenantId, EntityId entityId, String scope, List keys, TbCallback empty); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index 50541ccfcd..e3145b4398 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -218,12 +218,17 @@ public class TbEntityDataSubCtx { latestCtxValues.forEach((k, v) -> { TsValue update = latestUpdate.get(k); if (update != null) { - if (update.getTs() < v.getTs()) { - log.trace("[{}][{}][{}] Removed stale update for key: {} and ts: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k, update.getTs()); - latestUpdate.remove(k); - } else if ((update.getTs() == v.getTs() && update.getValue().equals(v.getValue()))) { - log.trace("[{}][{}][{}] Removed duplicate update for key: {} and ts: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k, update.getTs()); - latestUpdate.remove(k); + //Ignore notifications about deleted keys + if (!(update.getTs() == 0 && (update.getValue() == null || update.getValue().isEmpty()))) { + if (update.getTs() < v.getTs()) { + log.trace("[{}][{}][{}] Removed stale update for key: {} and ts: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k, update.getTs()); + latestUpdate.remove(k); + } else if ((update.getTs() == v.getTs() && update.getValue().equals(v.getValue()))) { + log.trace("[{}][{}][{}] Removed duplicate update for key: {} and ts: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k, update.getTs()); + latestUpdate.remove(k); + } + } else { + log.trace("[{}][{}][{}] Received deleted notification for: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k); } } }); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java index d0b4d713f9..68d8c5d45f 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java @@ -34,6 +34,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.KeyValueType; import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionMgrMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeSubscriptionProto; import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeUpdateProto; +import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeDeleteProto; import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionCloseProto; import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionKetStateProto; import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionProto; @@ -182,6 +183,22 @@ public class TbSubscriptionUtils { return ToCoreMsg.newBuilder().setToSubscriptionMgrMsg(msgBuilder.build()).build(); } + public static ToCoreMsg toAttributesDeleteProto(TenantId tenantId, EntityId entityId, String scope, List keys) { + TbAttributeDeleteProto.Builder builder = TbAttributeDeleteProto.newBuilder(); + builder.setEntityType(entityId.getEntityType().name()); + builder.setEntityIdMSB(entityId.getId().getMostSignificantBits()); + builder.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); + builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + builder.setScope(scope); + builder.addAllKeys(keys); + + SubscriptionMgrMsgProto.Builder msgBuilder = SubscriptionMgrMsgProto.newBuilder(); + msgBuilder.setAttrDelete(builder); + 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()); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 665a787e6d..b80309f757 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -133,6 +133,13 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope, attributes)); } + @Override + public void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List keys, FutureCallback callback) { + ListenableFuture> deleteFuture = attrService.removeAll(tenantId, entityId, scope, keys); + addMainCallback(deleteFuture, callback); + addWsCallback(deleteFuture, success -> onAttributesDelete(tenantId, entityId, scope, keys)); + } + @Override public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value, FutureCallback callback) { saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry(key, value) @@ -171,6 +178,20 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio } } + private void onAttributesDelete(TenantId tenantId, EntityId entityId, String scope, List keys) { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); + if (currentPartitions.contains(tpi)) { + if (subscriptionManagerService.isPresent()) { + subscriptionManagerService.get().onAttributesDelete(tenantId, entityId, scope, keys, TbCallback.EMPTY); + } else { + log.warn("Possible misconfiguration because subscriptionManagerService is null!"); + } + } else { + TransportProtos.ToCoreMsg toCoreMsg = TbSubscriptionUtils.toAttributesDeleteProto(tenantId, entityId, scope, keys); + clusterService.pushMsgToCore(tpi, entityId.getId(), toCoreMsg, null); + } + } + private void onTimeSeriesUpdate(TenantId tenantId, EntityId entityId, List ts) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); if (currentPartitions.contains(tpi)) { diff --git a/application/src/main/resources/logback.xml b/application/src/main/resources/logback.xml index 01aec9ceb9..02101a8eaa 100644 --- a/application/src/main/resources/logback.xml +++ b/application/src/main/resources/logback.xml @@ -30,6 +30,8 @@ + + diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto index d3e850b39c..ac0ce48203 100644 --- a/common/queue/src/main/proto/queue.proto +++ b/common/queue/src/main/proto/queue.proto @@ -295,6 +295,16 @@ message TbAttributeUpdateProto { repeated TsKvProto data = 7; } +message TbAttributeDeleteProto { + string entityType = 1; + int64 entityIdMSB = 2; + int64 entityIdLSB = 3; + int64 tenantIdMSB = 4; + int64 tenantIdLSB = 5; + string scope = 6; + repeated string keys = 7; +} + message TbTimeSeriesUpdateProto { string entityType = 1; int64 entityIdMSB = 2; @@ -340,6 +350,7 @@ message SubscriptionMgrMsgProto { TbSubscriptionCloseProto subClose = 3; TbTimeSeriesUpdateProto tsUpdate = 4; TbAttributeUpdateProto attrUpdate = 5; + TbAttributeDeleteProto attrDelete = 6; } message LocalSubscriptionServiceMsgProto { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java index 0737d2b4b3..5dd5f51032 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java @@ -398,7 +398,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { //TODO: fetch last level only. //TODO: fetch distinct records. String lvlFilter = getLvlFilter(entityFilter.getMaxLevel()); - String selectFields = "SELECT tenant_id, customer_id, id, type, name, label FROM " + entityType.name() + " WHERE id in ( SELECT entity_id"; + String selectFields = "SELECT tenant_id, customer_id, id, created_time, type, name, label FROM " + entityType.name() + " WHERE id in ( SELECT entity_id"; String from = getQueryTemplate(entityFilter.getDirection()); String whereFilter = " WHERE re.relation_type = :where_relation_type AND re.to_type = :where_entity_type"; diff --git a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java index a6ef3935b0..9a90119c05 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java @@ -24,7 +24,7 @@ import java.util.Arrays; @RunWith(ClasspathSuite.class) @ClassnameFilters({ - "org.thingsboard.server.dao.service.sql.*SqlTest" + "org.thingsboard.server.dao.service.sql.EntityServiceSqlTest" }) public class SqlDaoServiceTestSuite { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java index 1e3d334e22..20b2484f9e 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java @@ -61,6 +61,7 @@ import org.thingsboard.server.dao.attributes.AttributesService; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.Comparator; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; @@ -490,13 +491,16 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { List loadedIds = loadedEntities.stream().map(EntityData::getEntityId).collect(Collectors.toList()); List deviceIds = devices.stream().map(Device::getId).collect(Collectors.toList()); - + deviceIds.sort(Comparator.comparing(EntityId::getId)); + loadedIds.sort(Comparator.comparing(EntityId::getId)); Assert.assertEquals(deviceIds, loadedIds); List loadedNames = loadedEntities.stream().map(entityData -> entityData.getLatest().get(EntityKeyType.ENTITY_FIELD).get("name").getValue()).collect(Collectors.toList()); List deviceNames = devices.stream().map(Device::getName).collect(Collectors.toList()); + Collections.sort(loadedNames); + Collections.sort(deviceNames); Assert.assertEquals(deviceNames, loadedNames); sortOrder = new EntityDataSortOrder( @@ -560,8 +564,11 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { loadedEntities.addAll(data.getData()); } Assert.assertEquals(67, loadedEntities.size()); - List loadedTemperatures = loadedEntities.stream().map(entityData -> - entityData.getLatest().get(EntityKeyType.ATTRIBUTE).get("temperature").getValue()).collect(Collectors.toList()); + List loadedTemperatures = new ArrayList<>(); + for (Device device : devices) { + loadedTemperatures.add(loadedEntities.stream().filter(entityData -> entityData.getEntityId().equals(device.getId())).findFirst().orElse(null) + .getLatest().get(EntityKeyType.ATTRIBUTE).get("temperature").getValue()); + } List deviceTemperatures = temperatures.stream().map(aLong -> Long.toString(aLong)).collect(Collectors.toList()); Assert.assertEquals(deviceTemperatures, loadedTemperatures); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java index d2e19652a2..742dc5b507 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java @@ -44,4 +44,7 @@ public interface RuleEngineTelemetryService { void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value, FutureCallback callback); + void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List keys, FutureCallback callback); + + } From 9942231d73f7b0516181ab1b0dd87665b0dd1df6 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 26 Jun 2020 11:59:28 +0300 Subject: [PATCH 2/4] All tests restored --- .../java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java index 9a90119c05..a6ef3935b0 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java @@ -24,7 +24,7 @@ import java.util.Arrays; @RunWith(ClasspathSuite.class) @ClassnameFilters({ - "org.thingsboard.server.dao.service.sql.EntityServiceSqlTest" + "org.thingsboard.server.dao.service.sql.*SqlTest" }) public class SqlDaoServiceTestSuite { From 714b7bd7798bde87e15da2aeb8ae985363d58353 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 26 Jun 2020 13:49:54 +0300 Subject: [PATCH 3/4] Tests Fix --- .../BaseEntityQueryControllerTest.java | 3 +++ .../dao/service/BaseEntityServiceTest.java | 16 ++++++++++++---- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseEntityQueryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseEntityQueryControllerTest.java index 2763d14b46..75de7293f0 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseEntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseEntityQueryControllerTest.java @@ -102,6 +102,7 @@ public abstract class BaseEntityQueryControllerTest extends AbstractControllerTe device.setType("default"); device.setLabel("testLabel" + (int) (Math.random() * 1000)); devices.add(doPost("/api/device", device, Device.class)); + Thread.sleep(1); } DeviceTypeFilter filter = new DeviceTypeFilter(); filter.setDeviceType("default"); @@ -141,6 +142,7 @@ public abstract class BaseEntityQueryControllerTest extends AbstractControllerTe device.setType("default"); device.setLabel("testLabel" + (int) (Math.random() * 1000)); devices.add(doPost("/api/device", device, Device.class)); + Thread.sleep(1); } DeviceTypeFilter filter = new DeviceTypeFilter(); @@ -210,6 +212,7 @@ public abstract class BaseEntityQueryControllerTest extends AbstractControllerTe device.setType("default"); device.setLabel("testLabel" + (int) (Math.random() * 1000)); devices.add(doPost("/api/device?accessToken=" + name, device, Device.class)); + Thread.sleep(1); long temperature = (long) (Math.random() * 100); temperatures.add(temperature); if (temperature > 45) { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java index 20b2484f9e..2521a359af 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java @@ -88,7 +88,7 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { } @Test - public void testCountEntitiesByQuery() { + public void testCountEntitiesByQuery() throws InterruptedException { List devices = new ArrayList<>(); for (int i = 0; i < 97; i++) { Device device = new Device(); @@ -131,7 +131,7 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { } @Test - public void testCountHierarchicalEntitiesByQuery() { + public void testCountHierarchicalEntitiesByQuery() throws InterruptedException { List assets = new ArrayList<>(); List devices = new ArrayList<>(); createTestHierarchy(assets, devices, new ArrayList<>(), new ArrayList<>(), new ArrayList<>(), new ArrayList<>()); @@ -408,7 +408,7 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { deviceService.deleteDevicesByTenantId(tenantId); } - private void createTestHierarchy(List assets, List devices, List consumptions, List highConsumptions, List temperatures, List highTemperatures) { + private void createTestHierarchy(List assets, List devices, List consumptions, List highConsumptions, List temperatures, List highTemperatures) throws InterruptedException { for (int i = 0; i < 5; i++) { Asset asset = new Asset(); asset.setTenantId(tenantId); @@ -416,6 +416,8 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { asset.setType("type" + i); asset.setLabel("AssetLabel" + i); asset = assetService.saveAsset(asset); + //TO make sure devices have different created time + Thread.sleep(1); assets.add(asset); EntityRelation er = new EntityRelation(); er.setFrom(tenantId); @@ -435,6 +437,8 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { device.setType("default" + j); device.setLabel("testLabel" + (int) (Math.random() * 1000)); device = deviceService.saveDevice(device); + //TO make sure devices have different created time + Thread.sleep(1); devices.add(device); er = new EntityRelation(); er.setFrom(asset.getId()); @@ -452,7 +456,7 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { } @Test - public void testSimpleFindEntityDataByQuery() { + public void testSimpleFindEntityDataByQuery() throws InterruptedException { List devices = new ArrayList<>(); for (int i = 0; i < 97; i++) { Device device = new Device(); @@ -460,6 +464,8 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { device.setName("Device" + i); device.setType("default"); device.setLabel("testLabel" + (int) (Math.random() * 1000)); + //TO make sure devices have different created time + Thread.sleep(1); devices.add(deviceService.saveDevice(device)); } @@ -529,6 +535,8 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest { device.setType("default"); device.setLabel("testLabel" + (int) (Math.random() * 1000)); devices.add(deviceService.saveDevice(device)); + //TO make sure devices have different created time + Thread.sleep(1); long temperature = (long) (Math.random() * 100); temperatures.add(temperature); if (temperature > 45) { From 1b14f26dae5568c43f2b737022eca7bf15d2f652 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 26 Jun 2020 14:53:57 +0300 Subject: [PATCH 4/4] Tests Fix --- .../server/controller/BaseWebsocketApiTest.java | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java index 0a2e798b08..fee52548d3 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -455,22 +455,18 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getTs()); Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getValue()); - AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); - List tsData = Arrays.asList(dataPoint1); - sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); + wsClient.registerWaitForUpdate(); Thread.sleep(100); - cmd = new EntityDataCmd(1, edq, null, latestCmd, null); - wrapper = new TelemetryPluginCmdsWrapper(); - wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); + List tsData = Arrays.asList(dataPoint1); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); - wsClient.send(mapper.writeValueAsString(wrapper)); - msg = wsClient.waitForReply(); + msg = wsClient.waitForUpdate(); update = mapper.readValue(msg, EntityDataUpdate.class); Assert.assertEquals(1, update.getCmdId()); - List listData = update.getUpdate(); Assert.assertNotNull(listData); Assert.assertEquals(1, listData.size()); @@ -483,6 +479,7 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { AttributeKvEntry dataPoint2 = new BaseAttributeKvEntry(now, new LongDataEntry("serverAttributeKey", 52L)); wsClient.registerWaitForUpdate(); + Thread.sleep(100); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2)); msg = wsClient.waitForUpdate(); @@ -498,12 +495,14 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { //Sending update from the past, while latest value has new timestamp; wsClient.registerWaitForUpdate(); + Thread.sleep(100); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint1)); msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); Assert.assertNull(msg); //Sending duplicate update again wsClient.registerWaitForUpdate(); + Thread.sleep(100); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2)); msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); Assert.assertNull(msg);