Browse Source

fixed entity removal from CustomerData

pull/12580/head
dashevchenko 2 years ago
parent
commit
60736e1237
  1. 6
      application/src/main/java/org/thingsboard/server/service/edqs/EdqsSyncService.java
  2. 5
      application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java
  3. 18
      application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java
  4. 2
      common/data/src/main/java/org/thingsboard/server/common/data/ObjectType.java
  5. 2
      common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractQueryProcessor.java
  6. 5
      common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractRelationQueryProcessor.java
  7. 2
      common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/EntityListQueryProcessor.java
  8. 2
      common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/SingleEntityQueryProcessor.java
  9. 35
      common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java
  10. 2
      dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java
  11. 7
      dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java
  12. 8
      dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java

6
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));
}

5
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<Long, Long> condition) {
return await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> countByQuery(query),
result -> condition.test(result, expectedResult));
}
}

18
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<Device> 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<Long, Long> condition) throws Exception {
Long result = countByQuery(query);
assertThat(condition.test(result, expectedResult)).isTrue();
return result;
}
}

2
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() {

2
common/edqs/src/main/java/org/thingsboard/server/edqs/query/processor/AbstractQueryProcessor.java

@ -62,7 +62,7 @@ public abstract class AbstractQueryProcessor<T extends EntityFilter> 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));
}

5
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<T extends EntityFilter> 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<T extends EntityFilter> ext
var customerId = ctx.getCustomerId().getId();
List<SortableEntityData> result = new ArrayList<>();
for (EntityData<?> ed : entities) {
if (checkCustomer(customerId, ed)) {
if (checkCustomerId(customerId, ed)) {
result.add(toSortData(ed));
}
}

2
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<EntityData<?>> processor) {
processAll(ed -> {
if (checkCustomer(customerId, ed)) {
if (checkCustomerId(customerId, ed)) {
processor.accept(ed);
}
});

2
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<EntityData<?>> processor) {
processAll(ed -> {
if (checkCustomer(customerId, ed)) {
if (checkCustomerId(customerId, ed)) {
processor.accept(ed);
}
});

35
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();

2
dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java

@ -85,7 +85,7 @@ public class BaseQueueStatsService extends AbstractEntityService implements Queu
public PageData<QueueStats> 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

7
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<QueueStats> {
public interface QueueStatsDao extends Dao<QueueStats>, TenantEntityDao<QueueStats> {
QueueStats findByTenantIdQueueNameAndServiceId(TenantId tenantId, String queueName, String serviceId);
PageData<QueueStats> findByTenantId(TenantId tenantId, PageLink pageLink);
void deleteByTenantId(TenantId tenantId);
List<QueueStats> findByIds(TenantId tenantId, List<QueueStatsId> queueStatsIds);

8
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<QueueStatsEntity, QueueStat
}
@Override
public PageData<QueueStats> findByTenantId(TenantId tenantId, PageLink pageLink) {
public PageData<QueueStats> 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());

Loading…
Cancel
Save