diff --git a/application/src/main/java/org/thingsboard/server/service/edqs/EdqsSyncService.java b/application/src/main/java/org/thingsboard/server/service/edqs/EdqsSyncService.java index 04633892e5..c2750ae840 100644 --- a/application/src/main/java/org/thingsboard/server/service/edqs/EdqsSyncService.java +++ b/application/src/main/java/org/thingsboard/server/service/edqs/EdqsSyncService.java @@ -110,9 +110,9 @@ public abstract class EdqsSyncService { long ts = System.currentTimeMillis(); EntityType entityType = type.toEntityType(); Dao dao = entityDaoRegistry.getDao(entityType); - UUID lastFromEntityId = UUID.fromString("00000000-0000-0000-0000-000000000000"); + UUID lastId = UUID.fromString("00000000-0000-0000-0000-000000000000"); while (true) { - var batch = dao.findNextBatch(lastFromEntityId, 10000); + var batch = dao.findNextBatch(lastId, 10000); if (batch.isEmpty()) { break; } @@ -122,7 +122,7 @@ public abstract class EdqsSyncService { process(tenantId, type, new Entity(entityType, entityFields)); } EntityFields lastRecord = batch.get(batch.size() - 1); - lastFromEntityId = lastRecord.getId(); + lastId = lastRecord.getId(); } log.info("Finished synchronizing {} entities to EDQS in {} ms", type, (System.currentTimeMillis() - ts)); } diff --git a/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java index a75461faaa..baf9f98993 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java @@ -65,9 +65,4 @@ public class EdqsEntityQueryControllerTest extends EntityQueryControllerTest { result -> result == expectedResult); } - @Override - protected Long countByQueryAndCheck(EntityCountQuery query, long expectedResult, BiPredicate condition) { - return await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> countByQuery(query), - result -> condition.test(result, expectedResult)); - } } diff --git a/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java index 3941f16b88..508f66e584 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java @@ -81,7 +81,7 @@ import static org.awaitility.Awaitility.await; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @DaoSqlTest -public abstract class EntityQueryControllerTest extends AbstractControllerTest { +public class EntityQueryControllerTest extends AbstractControllerTest { private static final String CUSTOMER_USER_EMAIL = "entityQueryCustomer@thingsboard.org"; private static final String TENANT_PASSWORD = "testPassword1"; @@ -158,6 +158,13 @@ public abstract class EntityQueryControllerTest extends AbstractControllerTest { @Test public void testSysAdminCountEntitiesByQuery() throws Exception { + loginSysAdmin(); + + EntityTypeFilter allDeviceFilter = new EntityTypeFilter(); + allDeviceFilter.setEntityType(EntityType.DEVICE); + EntityCountQuery query = new EntityCountQuery(allDeviceFilter); + countByQueryAndCheck(query, 0); + loginTenantAdmin(); List devices = new ArrayList<>(); @@ -177,7 +184,7 @@ public abstract class EntityQueryControllerTest extends AbstractControllerTest { loginSysAdmin(); EntityCountQuery countQuery = new EntityCountQuery(filter); - countByQueryAndCheck(countQuery, 97, (actual, expected) -> actual >= expected); + countByQueryAndCheck(countQuery, 97); filter.setDeviceTypes(List.of("unknown")); countByQueryAndCheck(countQuery, 0); @@ -193,7 +200,7 @@ public abstract class EntityQueryControllerTest extends AbstractControllerTest { countQuery = new EntityCountQuery(entityListFilter); countByQueryAndCheck(countQuery, 97); - countByQueryAndCheck(countQuery, 97, (actual, expected) -> actual >= expected); + countByQueryAndCheck(countQuery, 97); } @Test @@ -906,9 +913,4 @@ public abstract class EntityQueryControllerTest extends AbstractControllerTest { return numericFilter; } - protected Long countByQueryAndCheck(EntityCountQuery query, long expectedResult, BiPredicate condition) throws Exception { - Long result = countByQuery(query); - assertThat(condition.test(result, expectedResult)).isTrue(); - return result; - } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ObjectType.java b/common/data/src/main/java/org/thingsboard/server/common/data/ObjectType.java index d3aa547bf0..bc5ec58213 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ObjectType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ObjectType.java @@ -73,7 +73,7 @@ public enum ObjectType { API_USAGE_STATE, ATTRIBUTE_KV, LATEST_TS_KV); static { - edqsTypes.addAll(Arrays.asList(TENANT, RELATION, ATTRIBUTE_KV, LATEST_TS_KV)); + edqsTypes.addAll(Arrays.asList(RELATION, ATTRIBUTE_KV, LATEST_TS_KV)); } public EntityType toEntityType() { diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractQueryProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractQueryProcessor.java index 29626dad1f..b795af08db 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractQueryProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractQueryProcessor.java @@ -62,7 +62,7 @@ public abstract class AbstractQueryProcessor implements } } - protected static boolean checkCustomer(UUID customerId, EntityData ed) { + protected static boolean checkCustomerId(UUID customerId, EntityData ed) { return customerId.equals(ed.getCustomerId()) || (ed.getEntityType() == EntityType.DASHBOARD && ed.getFields().getAssignedCustomerIds().contains(customerId)); } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractRelationQueryProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractRelationQueryProcessor.java index b5eaf26e6e..2a4e72ff0b 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractRelationQueryProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractRelationQueryProcessor.java @@ -16,7 +16,6 @@ package org.thingsboard.server.edqs.query.processor; import lombok.RequiredArgsConstructor; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.permission.QueryContext; import org.thingsboard.server.common.data.query.EntityFilter; import org.thingsboard.server.common.data.relation.EntitySearchDirection; @@ -79,7 +78,7 @@ public abstract class AbstractRelationQueryProcessor ext } else { var customerId = ctx.getCustomerId().getId(); for (EntityData ed : entities) { - if (checkCustomer(customerId, ed)) { + if (checkCustomerId(customerId, ed)) { result++; } } @@ -97,7 +96,7 @@ public abstract class AbstractRelationQueryProcessor ext var customerId = ctx.getCustomerId().getId(); List result = new ArrayList<>(); for (EntityData ed : entities) { - if (checkCustomer(customerId, ed)) { + if (checkCustomerId(customerId, ed)) { result.add(toSortData(ed)); } } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/EntityListQueryProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/EntityListQueryProcessor.java index 36016009bd..b319a4e13a 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/EntityListQueryProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/EntityListQueryProcessor.java @@ -41,7 +41,7 @@ public class EntityListQueryProcessor extends AbstractSingleEntityTypeQueryProce @Override protected void processCustomerQuery(UUID customerId, Consumer> processor) { processAll(ed -> { - if (checkCustomer(customerId, ed)) { + if (checkCustomerId(customerId, ed)) { processor.accept(ed); } }); diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/SingleEntityQueryProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/SingleEntityQueryProcessor.java index 19b32b5a93..55464b4529 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/SingleEntityQueryProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/SingleEntityQueryProcessor.java @@ -39,7 +39,7 @@ public class SingleEntityQueryProcessor extends AbstractSingleEntityTypeQueryPro @Override protected void processCustomerQuery(UUID customerId, Consumer> processor) { processAll(ed -> { - if (checkCustomer(customerId, ed)) { + if (checkCustomerId(customerId, ed)) { processor.accept(ed); } }); diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java index 7773c760f0..68e1c42589 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java @@ -195,24 +195,20 @@ public class TenantRepo { processFields(fields); entityData.setFields(entity.getFields()); - switch (entity.getType()) { - default -> { - UUID newCustomerId = fields.getCustomerId(); - UUID oldCustomerId = entityData.getCustomerId(); - entityData.setCustomerId(newCustomerId); - if (entityIdMismatch(oldCustomerId, newCustomerId)) { - if (oldCustomerId != null) { - CustomerData old = (CustomerData) getEntityMap(EntityType.CUSTOMER).get(oldCustomerId); - if (old != null) { - old.remove(entityData); - } - } - if (newCustomerId != null) { - CustomerData newData = (CustomerData) getEntityMap(EntityType.CUSTOMER).computeIfAbsent(newCustomerId, CustomerData::new); - newData.addOrUpdate(entityData); - } + UUID newCustomerId = fields.getCustomerId(); + UUID oldCustomerId = entityData.getCustomerId(); + entityData.setCustomerId(newCustomerId); + if (entityIdMismatch(oldCustomerId, newCustomerId)) { + if (oldCustomerId != null) { + CustomerData old = (CustomerData) getEntityMap(EntityType.CUSTOMER).get(oldCustomerId); + if (old != null) { + old.remove(entityData); } } + if (newCustomerId != null) { + CustomerData newData = (CustomerData) getEntityMap(EntityType.CUSTOMER).computeIfAbsent(newCustomerId, CustomerData::new); + newData.addOrUpdate(entityData); + } } } finally { entityUpdateLock.unlock(); @@ -228,6 +224,13 @@ public class TenantRepo { if (removed != null) { getEntitySet(entityType).remove(removed); edqsStatsService.ifPresent(statService -> statService.reportTenantEdqsObject(tenantId, ObjectType.fromEntityType(entityType), EdqsEventType.DELETED)); + UUID customerId = removed.getCustomerId(); + if (customerId != null) { + CustomerData customerData = (CustomerData) getEntityMap(EntityType.CUSTOMER).get(customerId); + if (customerData != null) { + customerData.remove(removed); + } + } } } finally { entityUpdateLock.unlock(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java index efc9dab41b..e3e6edad94 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java @@ -85,7 +85,7 @@ public class BaseQueueStatsService extends AbstractEntityService implements Queu public PageData findByTenantId(TenantId tenantId, PageLink pageLink) { log.trace("Executing findByTenantId, tenantId: [{}]", tenantId); Validator.validatePageLink(pageLink); - return queueStatsDao.findByTenantId(tenantId, pageLink); + return queueStatsDao.findAllByTenantId(tenantId, pageLink); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java b/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java index 5466a8afcc..018eb04d67 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java @@ -17,19 +17,16 @@ package org.thingsboard.server.dao.queue; import org.thingsboard.server.common.data.id.QueueStatsId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.queue.QueueStats; import org.thingsboard.server.dao.Dao; +import org.thingsboard.server.dao.TenantEntityDao; import java.util.List; -public interface QueueStatsDao extends Dao { +public interface QueueStatsDao extends Dao, TenantEntityDao { QueueStats findByTenantIdQueueNameAndServiceId(TenantId tenantId, String queueName, String serviceId); - PageData findByTenantId(TenantId tenantId, PageLink pageLink); - void deleteByTenantId(TenantId tenantId); List findByIds(TenantId tenantId, List queueStatsIds); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java index 97fedd693a..41e2de02a6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java @@ -21,6 +21,7 @@ import org.springframework.data.domain.Limit; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.edqs.fields.QueueStatsFields; import org.thingsboard.server.common.data.id.QueueStatsId; import org.thingsboard.server.common.data.id.TenantId; @@ -62,10 +63,15 @@ public class JpaQueueStatsDao extends JpaAbstractDao findByTenantId(TenantId tenantId, PageLink pageLink) { + public PageData findAllByTenantId(TenantId tenantId, PageLink pageLink) { return DaoUtil.toPageData(queueStatsRepository.findByTenantId(tenantId.getId(), pageLink.getTextSearch(), DaoUtil.toPageable(pageLink))); } + @Override + public ObjectType getType() { + return ObjectType.QUEUE_STATS; + } + @Override public void deleteByTenantId(TenantId tenantId) { queueStatsRepository.deleteByTenantId(tenantId.getId());