Browse Source

Merge pull request #14244 from dskarzh/fix/rule-node/customer-owned-entities

Improved support of customer-owned entities in `customer attributes` and `change originator` rule nodes
pull/14267/head
Viacheslav Klimov 9 months ago
committed by GitHub
parent
commit
1dfc5c0041
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityDaoService.java
  2. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityService.java
  3. 8
      common/dao-api/src/main/java/org/thingsboard/server/dao/entityview/EntityViewService.java
  4. 1
      common/data/src/main/java/org/thingsboard/server/common/data/HasCustomerId.java
  5. 7
      dao/src/main/java/org/thingsboard/server/dao/ai/AiModelServiceImpl.java
  6. 10
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
  7. 8
      dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java
  8. 11
      dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
  9. 8
      dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java
  10. 8
      dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java
  11. 9
      dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java
  12. 11
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java
  13. 26
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  14. 16
      dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java
  15. 10
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  16. 20
      dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java
  17. 30
      dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java
  18. 8
      dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java
  19. 15
      dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppBundleServiceImpl.java
  20. 14
      dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppServiceImpl.java
  21. 13
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java
  22. 14
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java
  23. 9
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java
  24. 14
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java
  25. 9
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java
  26. 9
      dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java
  27. 8
      dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java
  28. 36
      dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java
  29. 8
      dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java
  30. 34
      dao/src/main/java/org/thingsboard/server/dao/resource/BaseImageService.java
  31. 13
      dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java
  32. 35
      dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java
  33. 38
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  34. 9
      dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java
  35. 30
      dao/src/main/java/org/thingsboard/server/dao/tenant/TenantProfileServiceImpl.java
  36. 10
      dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java
  37. 11
      dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java
  38. 10
      dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java
  39. 12
      dao/src/main/java/org/thingsboard/server/dao/widget/WidgetTypeServiceImpl.java
  40. 48
      dao/src/main/java/org/thingsboard/server/dao/widget/WidgetsBundleServiceImpl.java
  41. 12
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDataNode.java
  42. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java
  43. 26
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java
  44. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  45. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java
  46. 14
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java
  47. 50
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java
  48. 2
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/TestDbCallbackExecutor.java
  49. 111
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java
  50. 90
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java
  51. 154
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java

3
common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityDaoService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.entity;
import com.google.common.util.concurrent.FluentFuture;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.HasId;
@ -26,6 +27,8 @@ public interface EntityDaoService {
Optional<HasId<?>> findEntity(TenantId tenantId, EntityId entityId);
FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId);
default long countByTenantId(TenantId tenantId) {
throw new IllegalArgumentException("Not implemented for " + getEntityType());
}

4
common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.entity;
import com.google.common.util.concurrent.FluentFuture;
import org.thingsboard.server.common.data.EntityInfo;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
@ -38,6 +39,8 @@ public interface EntityService {
Optional<CustomerId> fetchEntityCustomerId(TenantId tenantId, EntityId entityId);
FluentFuture<Optional<CustomerId>> fetchEntityCustomerIdAsync(TenantId tenantId, EntityId entityId);
Optional<HasId<?>> fetchEntity(TenantId tenantId, EntityId entityId);
Map<EntityId, EntityInfo> fetchEntityInfos(TenantId tenantId, CustomerId customerId, Set<EntityId> entityIds);
@ -47,4 +50,5 @@ public interface EntityService {
long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query);
PageData<EntityData> findEntityDataByQuery(TenantId tenantId, CustomerId customerId, EntityDataQuery query);
}

8
common/dao-api/src/main/java/org/thingsboard/server/dao/entityview/EntityViewService.java

@ -31,9 +31,6 @@ import org.thingsboard.server.dao.entity.EntityDaoService;
import java.util.List;
/**
* Created by Victor Basanets on 8/27/2017.
*/
public interface EntityViewService extends EntityDaoService {
EntityView saveEntityView(EntityView entityView);
@ -52,6 +49,8 @@ public interface EntityViewService extends EntityDaoService {
EntityView findEntityViewById(TenantId tenantId, EntityViewId entityViewId, boolean putInCache);
ListenableFuture<EntityView> findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId);
EntityView findEntityViewByTenantIdAndName(TenantId tenantId, String name);
ListenableFuture<EntityView> findEntityViewByTenantIdAndNameAsync(TenantId tenantId, String name);
@ -74,8 +73,6 @@ public interface EntityViewService extends EntityDaoService {
ListenableFuture<List<EntityView>> findEntityViewsByQuery(TenantId tenantId, EntityViewSearchQuery query);
ListenableFuture<EntityView> findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId);
ListenableFuture<List<EntityView>> findEntityViewsByTenantIdAndEntityIdAsync(TenantId tenantId, EntityId entityId);
List<EntityView> findEntityViewsByTenantIdAndEntityId(TenantId tenantId, EntityId entityId);
@ -95,4 +92,5 @@ public interface EntityViewService extends EntityDaoService {
PageData<EntityView> findEntityViewsByTenantIdAndEdgeId(TenantId tenantId, EdgeId edgeId, PageLink pageLink);
PageData<EntityView> findEntityViewsByTenantIdAndEdgeIdAndType(TenantId tenantId, EdgeId edgeId, String type, PageLink pageLink);
}

1
common/data/src/main/java/org/thingsboard/server/common/data/HasCustomerId.java

@ -20,4 +20,5 @@ import org.thingsboard.server.common.data.id.CustomerId;
public interface HasCustomerId {
CustomerId getCustomerId();
}

7
dao/src/main/java/org/thingsboard/server/dao/ai/AiModelServiceImpl.java

@ -36,6 +36,7 @@ import org.thingsboard.server.dao.sql.JpaExecutorService;
import java.util.Optional;
import java.util.Set;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@Service
@ -115,6 +116,12 @@ class AiModelServiceImpl extends CachedVersionedEntityService<AiModelCacheKey, A
.map(model -> model); // necessary to cast to HasId<?>
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return findAiModelByTenantIdAndIdAsync(tenantId, new AiModelId(entityId.getId()))
.transform(modelOpt -> modelOpt.map(model -> model), directExecutor()); // necessary to cast to HasId<?>
}
@Override
public long countByTenantId(TenantId tenantId) {
return aiModelDao.countByTenantId(tenantId);

10
dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.dao.alarm;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
@ -81,6 +80,7 @@ import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateEntityDataPageLink;
import static org.thingsboard.server.dao.service.Validator.validateId;
@ -97,8 +97,8 @@ public class BaseAlarmService extends AbstractCachedEntityService<TenantId, Page
private final AlarmDao alarmDao;
private final EntityService entityService;
@TransactionalEventListener(classes = AlarmTypesCacheEvictEvent.class)
@Override
@TransactionalEventListener
public void handleEvictEvent(AlarmTypesCacheEvictEvent event) {
TenantId tenantId = event.getTenantId();
cache.evict(tenantId);
@ -446,6 +446,12 @@ public class BaseAlarmService extends AbstractCachedEntityService<TenantId, Page
return Optional.ofNullable(findAlarmById(tenantId, new AlarmId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findAlarmByIdAsync(tenantId, new AlarmId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.ALARM;

8
dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.asset;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
@ -49,6 +50,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service("AssetProfileDaoService")
@ -323,6 +325,12 @@ public class AssetProfileServiceImpl extends CachedVersionedEntityService<AssetP
return Optional.ofNullable(findAssetProfileById(tenantId, new AssetProfileId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(assetProfileDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.ASSET_PROFILE;

11
dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java

@ -15,7 +15,7 @@
*/
package org.thingsboard.server.dao.asset;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
@ -62,6 +62,7 @@ import java.util.List;
import java.util.Optional;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.DaoUtil.toUUIDs;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validateIds;
@ -93,8 +94,8 @@ public class BaseAssetService extends AbstractCachedEntityService<AssetCacheKey,
@Autowired
private JpaExecutorService executor;
@TransactionalEventListener(classes = AssetCacheEvictEvent.class)
@Override
@TransactionalEventListener
public void handleEvictEvent(AssetCacheEvictEvent event) {
List<AssetCacheKey> keys = new ArrayList<>(2);
keys.add(new AssetCacheKey(event.getTenantId(), event.getNewName()));
@ -519,6 +520,12 @@ public class BaseAssetService extends AbstractCachedEntityService<AssetCacheKey,
return Optional.ofNullable(findAssetById(tenantId, new AssetId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findAssetByIdAsync(tenantId, new AssetId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public long countByTenantId(TenantId tenantId) {
return assetDao.countByTenantId(tenantId);

8
dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.cf;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@ -38,6 +39,7 @@ import org.thingsboard.server.dao.service.DataValidator;
import java.util.List;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@ -234,6 +236,12 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
return Optional.ofNullable(findById(tenantId, new CalculatedFieldId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(calculatedFieldDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.CALCULATED_FIELD;

8
dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.customer;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -55,6 +56,7 @@ import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service("CustomerDaoService")
@ -285,6 +287,12 @@ public class CustomerServiceImpl extends AbstractCachedEntityService<CustomerCac
return Optional.ofNullable(findCustomerById(tenantId, new CustomerId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findCustomerByIdAsync(tenantId, new CustomerId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public long countByTenantId(TenantId tenantId) {
return customerDao.countByTenantId(tenantId);

9
dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.dashboard;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@ -61,6 +62,7 @@ import java.util.List;
import java.util.Map;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service("DashboardDaoService")
@ -70,6 +72,7 @@ public class DashboardServiceImpl extends AbstractEntityService implements Dashb
public static final String INCORRECT_DASHBOARD_ID = "Incorrect dashboardId ";
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
@Autowired
private DashboardDao dashboardDao;
@ -424,6 +427,12 @@ public class DashboardServiceImpl extends AbstractEntityService implements Dashb
return Optional.ofNullable(findDashboardById(tenantId, new DashboardId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findDashboardByIdAsync(tenantId, new DashboardId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public long countByTenantId(TenantId tenantId) {
return dashboardDao.countByTenantId(tenantId);

11
dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.device;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -64,6 +65,7 @@ import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validateString;
@ -369,6 +371,12 @@ public class DeviceProfileServiceImpl extends CachedVersionedEntityService<Devic
return Optional.ofNullable(findDeviceProfileById(tenantId, new DeviceProfileId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(deviceProfileDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.DEVICE_PROFILE;
@ -432,8 +440,7 @@ public class DeviceProfileServiceImpl extends CachedVersionedEntityService<Devic
if (certificates.length > 1) {
return EncryptionUtil.certTrimNewLinesForChainInDeviceProfile(certificateValue);
}
} catch (CertificateException ignored) {
}
} catch (CertificateException ignored) {}
return EncryptionUtil.certTrimNewLines(certificateValue);
}

26
dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.device;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
@ -89,6 +90,7 @@ import java.util.List;
import java.util.Optional;
import java.util.UUID;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.DaoUtil.toUUIDs;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validateIds;
@ -126,23 +128,21 @@ public class DeviceServiceImpl extends CachedVersionedEntityService<DeviceCacheK
public Device findDeviceById(TenantId tenantId, DeviceId deviceId) {
log.trace("Executing findDeviceById [{}]", deviceId);
validateId(deviceId, id -> INCORRECT_DEVICE_ID + id);
if (TenantId.SYS_TENANT_ID.equals(tenantId)) {
return cache.get(new DeviceCacheKey(deviceId),
() -> deviceDao.findById(tenantId, deviceId.getId()));
} else {
return cache.get(new DeviceCacheKey(tenantId, deviceId),
() -> deviceDao.findDeviceByTenantIdAndId(tenantId, deviceId.getId()));
}
return findDeviceByIdInternal(tenantId, deviceId);
}
@Override
public ListenableFuture<Device> findDeviceByIdAsync(TenantId tenantId, DeviceId deviceId) {
log.trace("Executing findDeviceByIdAsync [{}]", deviceId);
validateId(deviceId, id -> INCORRECT_DEVICE_ID + id);
return executor.submit(() -> findDeviceByIdInternal(tenantId, deviceId));
}
private Device findDeviceByIdInternal(TenantId tenantId, DeviceId deviceId) {
if (TenantId.SYS_TENANT_ID.equals(tenantId)) {
return deviceDao.findByIdAsync(tenantId, deviceId.getId());
return cache.get(new DeviceCacheKey(deviceId), () -> deviceDao.findById(tenantId, deviceId.getId()));
} else {
return deviceDao.findDeviceByTenantIdAndIdAsync(tenantId, deviceId.getId());
return cache.get(new DeviceCacheKey(tenantId, deviceId), () -> deviceDao.findDeviceByTenantIdAndId(tenantId, deviceId.getId()));
}
}
@ -256,8 +256,8 @@ public class DeviceServiceImpl extends CachedVersionedEntityService<DeviceCacheK
}
}
@TransactionalEventListener(classes = DeviceCacheEvictEvent.class)
@Override
@TransactionalEventListener
public void handleEvictEvent(DeviceCacheEvictEvent event) {
List<DeviceCacheKey> toEvict = new ArrayList<>(3);
toEvict.add(new DeviceCacheKey(event.getTenantId(), event.getNewName()));
@ -729,6 +729,12 @@ public class DeviceServiceImpl extends CachedVersionedEntityService<DeviceCacheK
return Optional.ofNullable(findDeviceById(tenantId, new DeviceId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findDeviceByIdAsync(tenantId, new DeviceId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.DEVICE;

16
dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.domain;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@ -39,17 +40,16 @@ import org.thingsboard.server.dao.service.validator.DomainDataValidator;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Slf4j
@Service
public class DomainServiceImpl extends AbstractEntityService implements DomainService {
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
@Autowired
private OAuth2ClientDao oauth2ClientDao;
@Autowired
@ -66,8 +66,7 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entityId(savedDomain.getId()).entity(savedDomain).build());
return savedDomain;
} catch (Exception e) {
checkConstraintViolation(e,
Map.of("domain_name_key", "Domain with such name and scheme already exists!"));
checkConstraintViolation(e, "domain_name_key", "Domain with such name and scheme already exists!");
throw e;
}
}
@ -142,6 +141,12 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe
return Optional.ofNullable(findDomainById(tenantId, new DomainId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(domainDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
@Transactional
public void deleteEntity(TenantId tenantId, EntityId id, boolean force) {
@ -163,4 +168,5 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe
public EntityType getEntityType() {
return EntityType.DOMAIN;
}
}

10
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.edge;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
@ -83,6 +84,7 @@ import java.util.Optional;
import java.util.UUID;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.DaoUtil.toUUIDs;
import static org.thingsboard.server.dao.edge.BaseRelatedEdgesService.RELATED_EDGES_CACHE_ITEMS;
import static org.thingsboard.server.dao.service.Validator.validateId;
@ -137,8 +139,8 @@ public class EdgeServiceImpl extends AbstractCachedEntityService<EdgeCacheKey, E
@Value("${edges.state.persistToTelemetry:false}")
private boolean persistToTelemetry;
@TransactionalEventListener(classes = EdgeCacheEvictEvent.class)
@Override
@TransactionalEventListener
public void handleEvictEvent(EdgeCacheEvictEvent event) {
List<EdgeCacheKey> keys = new ArrayList<>(2);
keys.add(new EdgeCacheKey(event.getTenantId(), event.getNewName()));
@ -629,6 +631,12 @@ public class EdgeServiceImpl extends AbstractCachedEntityService<EdgeCacheKey, E
return Optional.ofNullable(findEdgeById(tenantId, new EdgeId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findEdgeByIdAsync(tenantId, new EdgeId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.EDGE;

20
dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.entity;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Lazy;
@ -67,15 +68,13 @@ import java.util.concurrent.ExecutionException;
import java.util.function.Function;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.common.data.id.EntityId.NULL_UUID;
import static org.thingsboard.server.common.data.query.EntityFilterType.ENTITY_NAME;
import static org.thingsboard.server.common.data.query.EntityFilterType.ENTITY_TYPE;
import static org.thingsboard.server.dao.service.Validator.validateEntityDataPageLink;
import static org.thingsboard.server.dao.service.Validator.validateId;
/**
* Created by ashvayka on 04.05.17.
*/
@Service
@Slf4j
public class BaseEntityService extends AbstractEntityService implements EntityService {
@ -93,7 +92,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
@Autowired
@Lazy
EntityServiceRegistry entityServiceRegistry;
private EntityServiceRegistry entityServiceRegistry;
@Autowired
private EdqsService edqsService;
@ -198,6 +197,11 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
return fetchAndConvert(tenantId, entityId, this::getCustomerId);
}
@Override
public FluentFuture<Optional<CustomerId>> fetchEntityCustomerIdAsync(TenantId tenantId, EntityId entityId) {
return fetchAndConvertAsync(tenantId, entityId, this::getCustomerId);
}
@Override
public Optional<NameLabelAndCustomerDetails> fetchNameLabelAndCustomerDetails(TenantId tenantId, EntityId entityId) {
log.trace("Executing fetchNameLabelAndCustomerDetails [{}]", entityId);
@ -239,6 +243,12 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
return entityOpt.map(converter);
}
private <T> FluentFuture<Optional<T>> fetchAndConvertAsync(TenantId tenantId, EntityId entityId, Function<HasId<?>, T> converter) {
EntityDaoService entityDaoService = entityServiceRegistry.getServiceByEntityType(entityId.getEntityType());
return entityDaoService.findEntityAsync(tenantId, entityId)
.transform(entityOpt -> entityOpt.map(converter), directExecutor());
}
private String getName(HasId<?> entity) {
return entity instanceof HasName ? ((HasName) entity).getName() : null;
}
@ -329,7 +339,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
}
if ((query.getEntityFields() == null || query.getEntityFields().isEmpty()) &&
(query.getLatestValues() == null || query.getLatestValues().isEmpty())) {
(query.getLatestValues() == null || query.getLatestValues().isEmpty())) {
return false;
}

30
dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.entityview;
import com.google.common.base.Function;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
@ -60,13 +61,11 @@ import java.util.List;
import java.util.Optional;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
import static org.thingsboard.server.dao.service.Validator.validateString;
/**
* Created by Victor Basanets on 8/28/2017.
*/
@Service("EntityViewDaoService")
@Slf4j
public class EntityViewServiceImpl extends CachedVersionedEntityService<EntityViewCacheKey, EntityViewCacheValue, EntityViewEvictEvent> implements EntityViewService {
@ -176,6 +175,17 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService<EntityVi
public EntityView findEntityViewById(TenantId tenantId, EntityViewId entityViewId, boolean putInCache) {
log.trace("Executing findEntityViewById [{}]", entityViewId);
validateId(entityViewId, id -> INCORRECT_ENTITY_VIEW_ID + id);
return findEntityViewByIdInternal(tenantId, entityViewId, putInCache);
}
@Override
public ListenableFuture<EntityView> findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId) {
log.trace("Executing findEntityViewByIdAsync [{}]", entityViewId);
validateId(entityViewId, id -> INCORRECT_ENTITY_VIEW_ID + id);
return service.submit(() -> findEntityViewByIdInternal(tenantId, entityViewId, true));
}
private EntityView findEntityViewByIdInternal(TenantId tenantId, EntityViewId entityViewId, boolean putInCache) {
EntityViewCacheValue value = cache.get(EntityViewCacheKey.byId(entityViewId), () -> {
EntityView entityView = entityViewDao.findById(tenantId, entityViewId.getId());
return new EntityViewCacheValue(entityView, null);
@ -190,7 +200,6 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService<EntityVi
return cache.getAndPutInTransaction(EntityViewCacheKey.byName(tenantId, name),
() -> entityViewDao.findEntityViewByTenantIdAndName(tenantId.getId(), name).orElse(null)
, EntityViewCacheValue::getEntityView, v -> new EntityViewCacheValue(v, null), true);
}
@Override
@ -307,13 +316,6 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService<EntityVi
return entityViews;
}
@Override
public ListenableFuture<EntityView> findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId) {
log.trace("Executing findEntityViewByIdAsync [{}]", entityViewId);
validateId(entityViewId, id -> INCORRECT_ENTITY_VIEW_ID + id);
return entityViewDao.findByIdAsync(tenantId, entityViewId.getId());
}
@Override
public ListenableFuture<List<EntityView>> findEntityViewsByTenantIdAndEntityIdAsync(TenantId tenantId, EntityId entityId) {
log.trace("Executing findEntityViewsByTenantIdAndEntityIdAsync, tenantId [{}], entityId [{}]", tenantId, entityId);
@ -479,6 +481,12 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService<EntityVi
return Optional.ofNullable(findEntityViewById(tenantId, new EntityViewId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findEntityViewByIdAsync(tenantId, new EntityViewId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.ENTITY_VIEW;

8
dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.job;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@ -44,6 +45,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.common.data.job.JobStatus.CANCELLED;
import static org.thingsboard.server.common.data.job.JobStatus.COMPLETED;
import static org.thingsboard.server.common.data.job.JobStatus.FAILED;
@ -241,6 +243,12 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
return Optional.ofNullable(findJobById(tenantId, (JobId) entityId));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(jobDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public void deleteEntity(TenantId tenantId, EntityId id, boolean force) {
jobDao.removeById(tenantId, id.getId());

15
dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppBundleServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.mobile;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@ -40,11 +41,11 @@ import org.thingsboard.server.dao.service.DataValidator;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.checkNotNull;
@Slf4j
@ -60,7 +61,6 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements
@Autowired
private DataValidator<MobileAppBundle> mobileAppBundleDataValidator;
@Override
public MobileAppBundle saveMobileAppBundle(TenantId tenantId, MobileAppBundle mobileAppBundle) {
log.trace("Executing saveMobileAppBundle [{}]", mobileAppBundle);
@ -71,8 +71,8 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements
return savedMobileApp;
} catch (Exception e) {
checkConstraintViolation(e,
Map.of("mobile_app_bundle_android_app_id_key", "Android mobile app is already configured in another bundle!",
"mobile_app_bundle_ios_app_id_key", "IOS mobile app is already configured in another bundle!"));
"mobile_app_bundle_android_app_id_key", "Android mobile app is already configured in another bundle!",
"mobile_app_bundle_ios_app_id_key", "IOS mobile app is already configured in another bundle!");
throw e;
}
}
@ -149,6 +149,12 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements
return Optional.ofNullable(findMobileAppBundleById(tenantId, new MobileAppBundleId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(mobileAppBundleDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
@Transactional
public void deleteEntity(TenantId tenantId, EntityId id, boolean force) {
@ -167,4 +173,5 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements
.collect(Collectors.toList());
mobileAppBundleInfo.setOauth2ClientInfos(clients);
}
}

14
dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.mobile;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@ -35,9 +36,10 @@ import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.Validator;
import java.util.Map;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Slf4j
@Service
public class MobileAppServiceImpl extends AbstractEntityService implements MobileAppService {
@ -58,8 +60,7 @@ public class MobileAppServiceImpl extends AbstractEntityService implements Mobil
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entity(savedMobileApp).build());
return savedMobileApp;
} catch (Exception e) {
checkConstraintViolation(e,
Map.of("mobile_app_pkg_name_platform_unq_key", "Mobile app with such package name and platform already exists!"));
checkConstraintViolation(e, "mobile_app_pkg_name_platform_unq_key", "Mobile app with such package name and platform already exists!");
throw e;
}
}
@ -88,6 +89,12 @@ public class MobileAppServiceImpl extends AbstractEntityService implements Mobil
return Optional.ofNullable(findMobileAppById(tenantId, new MobileAppId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(mobileAppDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
@Transactional
public void deleteEntity(TenantId tenantId, EntityId id, boolean force) {
@ -117,4 +124,5 @@ public class MobileAppServiceImpl extends AbstractEntityService implements Mobil
public EntityType getEntityType() {
return EntityType.MOBILE_APP;
}
}

13
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.notification;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.ApplicationEventPublisher;
@ -38,6 +39,8 @@ import org.thingsboard.server.dao.service.DataValidator;
import java.util.List;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Service
@Slf4j
@RequiredArgsConstructor
@ -129,13 +132,17 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
return Optional.ofNullable(findNotificationRequestById(tenantId, new NotificationRequestId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(notificationRequestDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.NOTIFICATION_REQUEST;
}
private static class NotificationRequestValidator extends DataValidator<NotificationRequest> {
}
private static class NotificationRequestValidator extends DataValidator<NotificationRequest> {}
}

14
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.notification;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.EntityType;
@ -33,9 +34,10 @@ import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent;
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Service
@RequiredArgsConstructor
public class DefaultNotificationRuleService extends AbstractEntityService implements NotificationRuleService, EntityDaoService {
@ -56,9 +58,7 @@ public class DefaultNotificationRuleService extends AbstractEntityService implem
.created(notificationRule.getId() == null).build());
return savedRule;
} catch (Exception e) {
checkConstraintViolation(e, Map.of(
"uq_notification_rule_name", "Notification rule with such name already exists"
));
checkConstraintViolation(e, "uq_notification_rule_name", "Notification rule with such name already exists");
throw e;
}
}
@ -114,6 +114,12 @@ public class DefaultNotificationRuleService extends AbstractEntityService implem
return Optional.ofNullable(findNotificationRuleById(tenantId, new NotificationRuleId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(notificationRuleDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.NOTIFICATION_RULE;

9
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.notification;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@ -37,6 +38,8 @@ import org.thingsboard.server.dao.sql.query.EntityKeyMapping;
import java.util.Optional;
import java.util.Set;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Service
@Slf4j
@RequiredArgsConstructor
@ -95,6 +98,12 @@ public class DefaultNotificationService implements NotificationService, EntityDa
return Optional.ofNullable(findNotificationById(tenantId, new NotificationId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(notificationDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.NOTIFICATION;

14
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.notification;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@ -47,13 +48,14 @@ import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent;
import org.thingsboard.server.dao.user.UserService;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.stream.Collectors;
import static org.apache.commons.collections4.CollectionUtils.isNotEmpty;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Service
@Slf4j
@RequiredArgsConstructor
@ -72,9 +74,7 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl
.created(notificationTarget.getId() == null).build());
return savedTarget;
} catch (Exception e) {
checkConstraintViolation(e, Map.of(
"uq_notification_target_name", "Recipients group with such name already exists"
));
checkConstraintViolation(e, "uq_notification_target_name", "Recipients group with such name already exists");
throw e;
}
}
@ -229,6 +229,12 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl
return Optional.ofNullable(findNotificationTargetById(tenantId, new NotificationTargetId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(notificationTargetDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.NOTIFICATION_TARGET;

9
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.notification;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.EntityType;
@ -37,6 +38,8 @@ import java.util.List;
import java.util.Map;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Service
@RequiredArgsConstructor
public class DefaultNotificationTemplateService extends AbstractEntityService implements NotificationTemplateService, EntityDaoService {
@ -144,6 +147,12 @@ public class DefaultNotificationTemplateService extends AbstractEntityService im
return Optional.ofNullable(findNotificationTemplateById(tenantId, new NotificationTemplateId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(notificationTemplateDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.NOTIFICATION_TEMPLATE;

9
dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.oauth2;
import com.google.common.util.concurrent.FluentFuture;
import jakarta.transaction.Transactional;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -41,6 +42,7 @@ import java.util.List;
import java.util.Optional;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Slf4j
@Service("OAuth2ClientService")
@ -107,7 +109,6 @@ public class OAuth2ClientServiceImpl extends AbstractEntityService implements OA
.tenantId(tenantId)
.entityId(oAuth2ClientId)
.build());
}
@Override
@ -149,6 +150,12 @@ public class OAuth2ClientServiceImpl extends AbstractEntityService implements OA
return Optional.ofNullable(findOAuth2ClientById(tenantId, new OAuth2ClientId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(oauth2ClientDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
@Transactional
public void deleteEntity(TenantId tenantId, EntityId id, boolean force) {

8
dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.ota;
import com.google.common.hash.HashFunction;
import com.google.common.hash.Hashing;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@ -47,6 +48,7 @@ import org.thingsboard.server.dao.service.PaginatedRemover;
import java.nio.ByteBuffer;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@ -254,6 +256,12 @@ public class BaseOtaPackageService extends AbstractCachedEntityService<OtaPackag
return Optional.ofNullable(findOtaPackageInfoById(tenantId, new OtaPackageId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findOtaPackageInfoByIdAsync(tenantId, new OtaPackageId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.OTA_PACKAGE;

36
dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.queue;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
@ -42,6 +43,9 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import java.util.List;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service("QueueDaoService")
@Slf4j
@RequiredArgsConstructor
@ -127,7 +131,7 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ
@Override
public void deleteQueuesByTenantId(TenantId tenantId) {
Validator.validateId(tenantId, "Incorrect tenant id for delete queues request.");
validateId(tenantId, __ -> "Incorrect tenant id for delete queues request.");
tenantQueuesRemover.removeEntities(tenantId, tenantId);
}
@ -141,24 +145,30 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ
return Optional.ofNullable(findQueueById(tenantId, new QueueId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(queueDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.QUEUE;
}
private PaginatedRemover<TenantId, Queue> tenantQueuesRemover =
new PaginatedRemover<>() {
private final PaginatedRemover<TenantId, Queue> tenantQueuesRemover = new PaginatedRemover<>() {
@Override
protected PageData<Queue> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return queueDao.findQueuesByTenantId(id, pageLink);
}
@Override
protected PageData<Queue> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return queueDao.findQueuesByTenantId(id, pageLink);
}
@Override
protected void removeEntity(TenantId tenantId, Queue entity) {
deleteQueue(tenantId, entity.getId());
}
};
@Override
protected void removeEntity(TenantId tenantId, Queue entity) {
deleteQueue(tenantId, entity.getId());
}
};
private TenantId getSystemOrIsolatedTenantId(TenantId tenantId) {
if (!tenantId.equals(TenantId.SYS_TENANT_ID)) {
@ -167,7 +177,7 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ
return tenantId;
}
}
return TenantId.SYS_TENANT_ID;
}
}

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

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.queue;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@ -35,6 +36,7 @@ import org.thingsboard.server.dao.service.Validator;
import java.util.List;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validateIds;
@ -106,6 +108,12 @@ public class BaseQueueStatsService extends AbstractEntityService implements Queu
return Optional.ofNullable(findQueueStatsById(tenantId, new QueueStatsId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(queueStatsDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.QUEUE_STATS;

34
dao/src/main/java/org/thingsboard/server/dao/resource/BaseImageService.java

@ -17,11 +17,11 @@ package org.thingsboard.server.dao.resource;
import com.fasterxml.jackson.databind.JsonNode;
import jakarta.annotation.PostConstruct;
import lombok.Data;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.RandomStringUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.Strings;
import org.apache.commons.lang3.exception.ExceptionUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.springframework.stereotype.Service;
@ -356,8 +356,8 @@ public class BaseImageService extends BaseResourceService implements ImageServic
imageName = imageName + type + " image";
UpdateResult result = convertToImageUrl(entity.getTenantId(), imageName, entity.getImage(), Collections.emptyMap());
entity.setImage(result.getValue());
return result.isUpdated();
entity.setImage(result.value());
return result.updated();
}
@Transactional(noRollbackFor = Exception.class) // we don't want transaction to rollback in case of an image processing failure
@ -373,8 +373,8 @@ public class BaseImageService extends BaseResourceService implements ImageServic
Map<String, String> imagesLinks = getResourcesLinks(widgetTypeDetails.getResources());
UpdateResult result = convertToImageUrl(tenantId, prefix + " image", widgetTypeDetails.getImage(), imagesLinks);
boolean updated = result.isUpdated();
widgetTypeDetails.setImage(result.getValue());
boolean updated = result.updated();
widgetTypeDetails.setImage(result.value());
if (widgetTypeDetails.getDescriptor().isObject()) {
JsonNode defaultConfig = widgetTypeDetails.getDefaultConfig();
@ -397,8 +397,8 @@ public class BaseImageService extends BaseResourceService implements ImageServic
Map<String, String> imagesLinks = getResourcesLinks(dashboard.getResources());
var result = convertToImageUrl(tenantId, prefix + " image", dashboard.getImage(), imagesLinks);
boolean updated = result.isUpdated();
dashboard.setImage(result.getValue());
boolean updated = result.updated();
dashboard.setImage(result.value());
updated |= convertToImageUrlsByMapping(tenantId, DASHBOARD_BASE64_MAPPING, Collections.singletonMap("prefix", prefix), dashboard.getConfiguration(), imagesLinks);
updated |= convertToImageUrls(tenantId, prefix, dashboard.getConfiguration(), imagesLinks);
@ -409,10 +409,10 @@ public class BaseImageService extends BaseResourceService implements ImageServic
AtomicBoolean updated = new AtomicBoolean(false);
JacksonUtil.replaceAllByMapping(configuration, mapping, templateParams, (name, value) -> {
UpdateResult result = convertToImageUrl(tenantId, name, value, links);
if (result.isUpdated()) {
if (result.updated()) {
updated.set(true);
}
return result.getValue();
return result.value();
});
return updated.get();
}
@ -514,10 +514,10 @@ public class BaseImageService extends BaseResourceService implements ImageServic
AtomicBoolean updated = new AtomicBoolean(false);
JacksonUtil.replaceAll(root, title, (path, value) -> {
UpdateResult result = convertToImageUrl(tenantId, path, value, true, links);
if (result.isUpdated()) {
if (result.updated()) {
updated.set(true);
}
return result.getValue();
return result.value();
});
return updated.get();
}
@ -678,16 +678,18 @@ public class BaseImageService extends BaseResourceService implements ImageServic
private String getImageLink(String value) {
if (value.startsWith(DataConstants.TB_IMAGE_PREFIX + "/api/images")) {
return StringUtils.removeStart(value, DataConstants.TB_IMAGE_PREFIX);
return Strings.CS.removeStart(value, DataConstants.TB_IMAGE_PREFIX);
} else {
return null;
}
}
@Data(staticConstructor = "of")
private static class UpdateResult {
private final boolean updated;
private final String value;
private record UpdateResult(boolean updated, String value) {
static UpdateResult of(boolean updated, String value) {
return new UpdateResult(updated, value);
}
}
}

13
dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.resource;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.TextNode;
import com.google.common.hash.Hashing;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import jakarta.annotation.PostConstruct;
import lombok.RequiredArgsConstructor;
@ -79,6 +80,7 @@ import java.util.UUID;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.UnaryOperator;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.common.data.StringUtils.isNotEmpty;
import static org.thingsboard.server.dao.device.DeviceServiceImpl.INCORRECT_TENANT_ID;
import static org.thingsboard.server.dao.service.Validator.validateId;
@ -107,7 +109,8 @@ public class BaseResourceService extends AbstractCachedEntityService<ResourceInf
generalResourceContainerDaoMap.put(EntityType.RULE_CHAIN, ruleChainDao);
}
@Autowired @Lazy
@Autowired
@Lazy
private ImageService imageService;
private static final Map<String, String> DASHBOARD_RESOURCES_MAPPING = Map.of(
@ -275,7 +278,7 @@ public class BaseResourceService extends AbstractCachedEntityService<ResourceInf
@Override
public TbResource toResource(TenantId tenantId, ResourceExportData exportData) {
if (exportData.getType() == ResourceType.IMAGE || exportData.getSubType() == ResourceSubType.IMAGE
|| exportData.getSubType() == ResourceSubType.SCADA_SYMBOL) {
|| exportData.getSubType() == ResourceSubType.SCADA_SYMBOL) {
throw new IllegalArgumentException("Image import not supported");
}
@ -457,6 +460,12 @@ public class BaseResourceService extends AbstractCachedEntityService<ResourceInf
return Optional.ofNullable(findResourceInfoById(tenantId, new TbResourceId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findResourceInfoByIdAsync(tenantId, new TbResourceId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.TB_RESOURCE;

35
dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.rpc;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@ -33,6 +34,7 @@ import org.thingsboard.server.dao.service.PaginatedRemover;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@ -40,6 +42,7 @@ import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@Slf4j
@RequiredArgsConstructor
public class BaseRpcService implements RpcService {
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
public static final String INCORRECT_RPC_ID = "Incorrect rpcId ";
@ -113,21 +116,29 @@ public class BaseRpcService implements RpcService {
return Optional.ofNullable(findById(tenantId, new RpcId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findRpcByIdAsync(tenantId, new RpcId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.RPC;
}
private PaginatedRemover<TenantId, Rpc> tenantRpcRemover =
new PaginatedRemover<>() {
@Override
protected PageData<Rpc> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return rpcDao.findAllRpcByTenantId(id, pageLink);
}
@Override
protected void removeEntity(TenantId tenantId, Rpc entity) {
deleteRpc(tenantId, entity.getId());
}
};
private final PaginatedRemover<TenantId, Rpc> tenantRpcRemover = new PaginatedRemover<>() {
@Override
protected PageData<Rpc> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return rpcDao.findAllRpcByTenantId(id, pageLink);
}
@Override
protected void removeEntity(TenantId tenantId, Rpc entity) {
deleteRpc(tenantId, entity.getId());
}
};
}

38
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.rule;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections4.CollectionUtils;
@ -28,6 +29,7 @@ import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.BaseDataWithAdditionalInfo;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.edge.Edge;
@ -38,6 +40,7 @@ import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.data.page.PageLink;
@ -80,6 +83,7 @@ import java.util.UUID;
import java.util.function.Function;
import java.util.stream.Collectors;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.common.data.DataConstants.TENANT;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validateIds;
@ -87,9 +91,6 @@ import static org.thingsboard.server.dao.service.Validator.validatePageLink;
import static org.thingsboard.server.dao.service.Validator.validatePositiveNumber;
import static org.thingsboard.server.dao.service.Validator.validateString;
/**
* Created by igor on 3/12/18.
*/
@Service("RuleChainDaoService")
@Slf4j
public class BaseRuleChainService extends AbstractEntityService implements RuleChainService {
@ -98,6 +99,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
public static final String TB_RULE_CHAIN_INPUT_NODE = "org.thingsboard.rule.engine.flow.TbRuleChainInputNode";
@Autowired
private RuleChainDao ruleChainDao;
@ -260,7 +262,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
firstRuleNodeId = nodes.get(ruleChainMetaData.getFirstNodeIndex()).getId();
}
if ((ruleChain.getFirstRuleNodeId() != null && !ruleChain.getFirstRuleNodeId().equals(firstRuleNodeId))
|| (ruleChain.getFirstRuleNodeId() == null && firstRuleNodeId != null)) {
|| (ruleChain.getFirstRuleNodeId() == null && firstRuleNodeId != null)) {
ruleChain.setFirstRuleNodeId(firstRuleNodeId);
}
if (ruleChainMetaData.getConnections() != null) {
@ -876,6 +878,17 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
return Optional.ofNullable(hasId);
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
ListenableFuture<? extends BaseDataWithAdditionalInfo<? extends UUIDBased>> future;
if (entityId.getEntityType() == EntityType.RULE_NODE) {
future = findRuleNodeByIdAsync(tenantId, new RuleNodeId(entityId.getId()));
} else {
future = findRuleChainByIdAsync(tenantId, new RuleChainId(entityId.getId()));
}
return FluentFuture.from(future).transform(Optional::ofNullable, directExecutor());
}
@Override
public long countByTenantId(TenantId tenantId) {
return ruleChainDao.countByTenantId(tenantId);
@ -919,18 +932,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
ComponentClusteringMode nodeConfigType = ReflectionUtils.getAnnotationProperty(ruleNode.getType(),
"org.thingsboard.rule.engine.api.RuleNode", "clusteringMode");
switch (nodeConfigType) {
case ENABLED:
singletonMode = false;
break;
case SINGLETON:
singletonMode = true;
break;
case USER_PREFERENCE:
default:
singletonMode = ruleNode.isSingletonMode();
break;
}
singletonMode = switch (nodeConfigType) {
case ENABLED -> false;
case SINGLETON -> true;
default -> ruleNode.isSingletonMode();
};
} catch (Exception e) {
log.warn("Failed to get clustering mode: {}", ExceptionUtils.getRootCauseMessage(e));
singletonMode = false;

9
dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.settings;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@ -31,6 +32,8 @@ import org.thingsboard.server.dao.service.Validator;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
@Service
@Slf4j
public class AdminSettingsServiceImpl implements AdminSettingsService {
@ -106,6 +109,12 @@ public class AdminSettingsServiceImpl implements AdminSettingsService {
return Optional.ofNullable(adminSettingsDao.findById(tenantId, entityId.getId()));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(adminSettingsDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.ADMIN_SETTINGS;

30
dao/src/main/java/org/thingsboard/server/dao/tenant/TenantProfileServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.tenant;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
@ -44,6 +45,7 @@ import java.util.List;
import java.util.Optional;
import java.util.UUID;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.common.util.DebugModeUtil.DEBUG_MODE_DEFAULT_DURATION_MINUTES;
import static org.thingsboard.server.dao.service.Validator.validateId;
@ -230,23 +232,29 @@ public class TenantProfileServiceImpl extends AbstractCachedEntityService<Tenant
return Optional.ofNullable(findTenantProfileById(tenantId, new TenantProfileId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(tenantProfileDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.TENANT_PROFILE;
}
private final PaginatedRemover<String, TenantProfile> tenantProfilesRemover =
new PaginatedRemover<>() {
private final PaginatedRemover<String, TenantProfile> tenantProfilesRemover = new PaginatedRemover<>() {
@Override
protected PageData<TenantProfile> findEntities(TenantId tenantId, String id, PageLink pageLink) {
return tenantProfileDao.findTenantProfiles(tenantId, pageLink);
}
@Override
protected PageData<TenantProfile> findEntities(TenantId tenantId, String id, PageLink pageLink) {
return tenantProfileDao.findTenantProfiles(tenantId, pageLink);
}
@Override
protected void removeEntity(TenantId tenantId, TenantProfile entity) {
removeTenantProfile(tenantId, entity, entity.isDefault());
}
@Override
protected void removeEntity(TenantId tenantId, TenantProfile entity) {
removeTenantProfile(tenantId, entity, entity.isDefault());
}
};
};
}

10
dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.tenant;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -51,6 +52,7 @@ import java.util.List;
import java.util.Optional;
import java.util.function.Consumer;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service("TenantDaoService")
@ -85,8 +87,8 @@ public class TenantServiceImpl extends AbstractCachedEntityService<TenantId, Ten
@Autowired
protected TbTransactionalCache<TenantId, Boolean> existsTenantCache;
@TransactionalEventListener(classes = TenantEvictEvent.class)
@Override
@TransactionalEventListener
public void handleEvictEvent(TenantEvictEvent event) {
TenantId tenantId = event.getTenantId();
cache.evict(tenantId);
@ -245,6 +247,12 @@ public class TenantServiceImpl extends AbstractCachedEntityService<TenantId, Ten
return Optional.ofNullable(findTenantById(TenantId.fromUUID(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findTenantByIdAsync(tenantId, TenantId.fromUUID(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.TENANT;

11
dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.usagerecord;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
@ -48,11 +49,13 @@ import java.util.List;
import java.util.Objects;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service("ApiUsageStateDaoService")
@Slf4j
public class ApiUsageStateServiceImpl extends AbstractEntityService implements ApiUsageStateService {
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
private final ApiUsageStateDao apiUsageStateDao;
@ -161,7 +164,7 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A
public ApiUsageState update(ApiUsageState apiUsageState) {
log.trace("Executing save [{}]", apiUsageState.getTenantId());
validateId(apiUsageState.getTenantId(), id -> INCORRECT_TENANT_ID + id);
validateId(apiUsageState.getId(), "Can't save new usage state. Only update is allowed!");
validateId(apiUsageState.getId(), __ -> "Can't save new usage state. Only update is allowed!");
apiUsageState.setVersion(null);
ApiUsageState savedState = apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState);
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(savedState.getTenantId()).entityId(savedState.getId())
@ -195,6 +198,12 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A
return Optional.ofNullable(findApiUsageStateById(tenantId, new ApiUsageStateId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(apiUsageStateDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Transactional
@Override
public void deleteByTenantId(TenantId tenantId) {

10
dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.user;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@ -72,6 +73,7 @@ import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.common.data.StringUtils.generateSafeToken;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@ -103,8 +105,8 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
private final EntityCountService countService;
private final JpaExecutorService executor;
@TransactionalEventListener(classes = UserCacheEvictEvent.class)
@Override
@TransactionalEventListener
public void handleEvictEvent(UserCacheEvictEvent event) {
List<UserCacheKey> keys = new ArrayList<>(2);
keys.add(new UserCacheKey(event.tenantId(), event.newEmail()));
@ -568,6 +570,12 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
return Optional.ofNullable(findUserById(tenantId, new UserId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(findUserByIdAsync(tenantId, new UserId(entityId.getId())))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public long countByTenantId(TenantId tenantId) {
return userDao.countByTenantId(tenantId);

12
dao/src/main/java/org/thingsboard/server/dao/widget/WidgetTypeServiceImpl.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.widget;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections4.CollectionUtils;
import org.springframework.beans.factory.annotation.Autowired;
@ -47,6 +48,7 @@ import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.service.Validator.validateIds;
@Service("WidgetTypeDaoService")
@ -54,8 +56,6 @@ import static org.thingsboard.server.dao.service.Validator.validateIds;
public class WidgetTypeServiceImpl implements WidgetTypeService {
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
public static final String INCORRECT_RESOURCE_ID = "Incorrect resourceId ";
public static final String INCORRECT_BUNDLE_ALIAS = "Incorrect bundleAlias ";
public static final String INCORRECT_WIDGETS_BUNDLE_ID = "Incorrect widgetsBundleId ";
@Autowired
@ -281,6 +281,12 @@ public class WidgetTypeServiceImpl implements WidgetTypeService {
return Optional.ofNullable(findWidgetTypeById(tenantId, new WidgetTypeId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(widgetTypeDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.WIDGET_TYPE;
@ -303,6 +309,7 @@ public class WidgetTypeServiceImpl implements WidgetTypeService {
protected void removeEntity(TenantId tenantId, WidgetTypeInfo entity) {
deleteWidgetType(tenantId, new WidgetTypeId(entity.getUuidId()));
}
};
private final PaginatedRemover<WidgetsBundleId, WidgetTypeInfo> bundleWidgetTypesRemover = new PaginatedRemover<>() {
@ -316,6 +323,7 @@ public class WidgetTypeServiceImpl implements WidgetTypeService {
protected void removeEntity(TenantId tenantId, WidgetTypeInfo widgetTypeInfo) {
deleteWidgetType(tenantId, widgetTypeInfo.getId());
}
};
}

48
dao/src/main/java/org/thingsboard/server/dao/widget/WidgetsBundleServiceImpl.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.widget;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher;
@ -33,12 +34,10 @@ import org.thingsboard.server.common.data.widget.WidgetType;
import org.thingsboard.server.common.data.widget.WidgetTypeDetails;
import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.common.data.widget.WidgetsBundleFilter;
import org.thingsboard.server.dao.entity.AbstractCachedEntityService;
import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent;
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.resource.ImageService;
import org.thingsboard.server.dao.resource.ResourceService;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.PaginatedRemover;
import org.thingsboard.server.dao.service.Validator;
@ -48,13 +47,15 @@ import java.util.List;
import java.util.Optional;
import java.util.stream.Stream;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static org.thingsboard.server.dao.entity.AbstractEntityService.checkConstraintViolation;
@Service("WidgetsBundleDaoService")
@Slf4j
public class WidgetsBundleServiceImpl implements WidgetsBundleService {
private static final int DEFAULT_WIDGETS_BUNDLE_LIMIT = 300;
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
public static final String INCORRECT_PAGE_LINK = "Incorrect page link ";
@Autowired
private WidgetsBundleDao widgetsBundleDao;
@ -71,9 +72,6 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService {
@Autowired
private ImageService imageService;
@Autowired
private ResourceService resourceService;
@Override
public WidgetsBundle findWidgetsBundleById(TenantId tenantId, WidgetsBundleId widgetsBundleId) {
log.trace("Executing findWidgetsBundleById [{}]", widgetsBundleId);
@ -92,9 +90,9 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService {
.entityId(result.getId()).created(widgetsBundle.getId() == null).build());
return result;
} catch (Exception e) {
AbstractCachedEntityService.checkConstraintViolation(e,
"uq_widgets_bundle_alias", "Widgets Bundle with such alias already exists!");
AbstractCachedEntityService.checkConstraintViolation(e, "widgets_bundle_external_id_unq_key", "Widgets Bundle with such external id already exists!");
checkConstraintViolation(e,
"uq_widgets_bundle_alias", "Widgets Bundle with such alias already exists!",
"widgets_bundle_external_id_unq_key", "Widgets Bundle with such external id already exists!");
throw e;
}
}
@ -248,9 +246,7 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService {
}
if (widgetsBundleDescriptor.has("widgetTypeFqns")) {
JsonNode widgetFqnsArrayJson = widgetsBundleDescriptor.get("widgetTypeFqns");
widgetFqnsArrayJson.forEach(fqnJson -> {
widgetTypeFqns.add(fqnJson.asText());
});
widgetFqnsArrayJson.forEach(fqnJson -> widgetTypeFqns.add(fqnJson.asText()));
}
widgetTypeService.updateWidgetsBundleWidgetFqns(TenantId.SYS_TENANT_ID, widgetsBundle.getId(), widgetTypeFqns);
});
@ -274,23 +270,29 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService {
return Optional.ofNullable(findWidgetsBundleById(tenantId, new WidgetsBundleId(entityId.getId())));
}
@Override
public FluentFuture<Optional<HasId<?>>> findEntityAsync(TenantId tenantId, EntityId entityId) {
return FluentFuture.from(widgetsBundleDao.findByIdAsync(tenantId, entityId.getId()))
.transform(Optional::ofNullable, directExecutor());
}
@Override
public EntityType getEntityType() {
return EntityType.WIDGETS_BUNDLE;
}
private PaginatedRemover<TenantId, WidgetsBundle> tenantWidgetsBundleRemover =
new PaginatedRemover<TenantId, WidgetsBundle>() {
private final PaginatedRemover<TenantId, WidgetsBundle> tenantWidgetsBundleRemover = new PaginatedRemover<>() {
@Override
protected PageData<WidgetsBundle> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return widgetsBundleDao.findTenantWidgetsBundlesByTenantId(id.getId(), pageLink);
}
@Override
protected PageData<WidgetsBundle> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return widgetsBundleDao.findTenantWidgetsBundlesByTenantId(id.getId(), pageLink);
}
@Override
protected void removeEntity(TenantId tenantId, WidgetsBundle entity) {
deleteWidgetsBundle(tenantId, new WidgetsBundleId(entity.getUuidId()));
}
@Override
protected void removeEntity(TenantId tenantId, WidgetsBundle entity) {
deleteWidgetsBundle(tenantId, new WidgetsBundleId(entity.getUuidId()));
}
};
};
}

12
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDataNode.java

@ -58,15 +58,9 @@ public abstract class TbAbstractGetEntityDataNode<T extends EntityId> extends Tb
protected void processDataAndTell(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) {
DataToFetch dataToFetch = config.getDataToFetch();
switch (dataToFetch) {
case ATTRIBUTES:
processAttributesKvEntryData(ctx, msg, entityId, msgDataAsJsonNode);
break;
case LATEST_TELEMETRY:
processTsKvEntryData(ctx, msg, entityId, msgDataAsJsonNode);
break;
case FIELDS:
processFieldsData(ctx, msg, entityId, msgDataAsJsonNode, true);
break;
case ATTRIBUTES -> processAttributesKvEntryData(ctx, msg, entityId, msgDataAsJsonNode);
case LATEST_TELEMETRY -> processTsKvEntryData(ctx, msg, entityId, msgDataAsJsonNode);
case FIELDS -> processFieldsData(ctx, msg, entityId, msgDataAsJsonNode, true);
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java

@ -17,8 +17,7 @@ package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.AsyncFunction;
import com.google.common.util.concurrent.Futures;
import com.google.common.base.Function;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbContext;
@ -54,12 +53,12 @@ public abstract class TbAbstractNodeWithFetchTo<C extends TbAbstractFetchToNodeC
protected abstract C loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException;
protected <I extends EntityId> AsyncFunction<I, I> checkIfEntityIsPresentOrThrow(String message) {
protected <I extends EntityId> Function<I, I> checkIfEntityIsPresentOrThrow(String message) {
return id -> {
if (id == null || id.isNullUid()) {
return Futures.immediateFailedFuture(new NoSuchElementException(message));
throw new NoSuchElementException(message);
}
return Futures.immediateFuture(id);
return id;
};
}

26
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java

@ -16,19 +16,22 @@
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.util.TbPair;
import java.util.NoSuchElementException;
import static com.google.common.util.concurrent.Futures.immediateFuture;
@RuleNode(
type = ComponentType.ENRICHMENT,
name = "customer attributes",
@ -44,8 +47,6 @@ import org.thingsboard.server.common.data.util.TbPair;
)
public class TbGetCustomerAttributeNode extends TbAbstractGetEntityDataNode<CustomerId> {
private static final String CUSTOMER_NOT_FOUND_MESSAGE = "Failed to find customer for entity with id: %s and type: %s";
@Override
protected TbGetEntityDataNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
var config = TbNodeUtils.convert(configuration, TbGetEntityDataNodeConfiguration.class);
@ -56,10 +57,19 @@ public class TbGetCustomerAttributeNode extends TbAbstractGetEntityDataNode<Cust
@Override
protected ListenableFuture<CustomerId> findEntityAsync(TbContext ctx, EntityId originator) {
return Futures.transformAsync(EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator),
checkIfEntityIsPresentOrThrow(String.format(CUSTOMER_NOT_FOUND_MESSAGE, originator.getId(), originator.getEntityType().getNormalName())),
ctx.getDbCallbackExecutor()
);
if (originator.getEntityType() == EntityType.CUSTOMER) {
return immediateFuture((CustomerId) originator);
}
return ctx.getEntityService().fetchEntityCustomerIdAsync(ctx.getTenantId(), originator)
.transform(customerIdOpt -> {
if (customerIdOpt.isEmpty()) {
throw new NoSuchElementException("Originator not found");
}
if (customerIdOpt.get().isNullUid()) {
throw new IllegalStateException("Originator is not assigned to any customer");
}
return customerIdOpt.get();
}, ctx.getDbCallbackExecutor());
}
@Override

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java

@ -55,7 +55,7 @@ public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode<TbGetDevice
@Override
protected ListenableFuture<DeviceId> findEntityIdAsync(TbContext ctx, TbMsg msg) {
return Futures.transformAsync(
return Futures.transform(
EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery()),
checkIfEntityIsPresentOrThrow(RELATED_DEVICE_NOT_FOUND_MESSAGE),
ctx.getDbCallbackExecutor());

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java

@ -58,7 +58,7 @@ public class TbGetRelatedAttributeNode extends TbAbstractGetEntityDataNode<Entit
@Override
public ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
var relatedAttrConfig = (TbGetRelatedDataNodeConfiguration) config;
return Futures.transformAsync(
return Futures.transform(
EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery()),
checkIfEntityIsPresentOrThrow(RELATED_ENTITY_NOT_FOUND_MESSAGE),
ctx.getDbCallbackExecutor());

14
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java

@ -25,7 +25,6 @@ import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.util.EntitiesAlarmOriginatorIdAsyncLoader;
import org.thingsboard.rule.engine.util.EntitiesByNameAndTypeLoader;
import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader;
import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
@ -36,6 +35,7 @@ import org.thingsboard.server.common.msg.TbMsg;
import java.util.List;
import java.util.NoSuchElementException;
import static com.google.common.util.concurrent.Futures.immediateFuture;
import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY;
import static org.thingsboard.rule.engine.transform.OriginatorSource.RELATED;
@ -74,16 +74,20 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode<TbChangeOrig
if (newOriginator == null || newOriginator.isNullUid()) {
return Futures.immediateFailedFuture(new NoSuchElementException("Failed to find new originator!"));
}
return Futures.immediateFuture(List.of(ctx.transformMsgOriginator(msg, newOriginator)));
return immediateFuture(List.of(ctx.transformMsgOriginator(msg, newOriginator)));
}, ctx.getDbCallbackExecutor());
}
private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, TbMsg msg) {
switch (config.getOriginatorSource()) {
case CUSTOMER:
return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator());
if (msg.getOriginator().getEntityType() == EntityType.CUSTOMER) {
return immediateFuture(msg.getOriginator());
}
return ctx.getEntityService().fetchEntityCustomerIdAsync(ctx.getTenantId(), msg.getOriginator())
.transform(customerIdOpt -> customerIdOpt.orElse(null), ctx.getDbCallbackExecutor());
case TENANT:
return Futures.immediateFuture(ctx.getTenantId());
return immediateFuture(ctx.getTenantId());
case RELATED:
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, msg.getOriginator(), config.getRelationsQuery());
case ALARM_ORIGINATOR:
@ -93,7 +97,7 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode<TbChangeOrig
String entityName = TbNodeUtils.processPattern(config.getEntityNamePattern(), msg);
try {
EntityId targetEntity = EntitiesByNameAndTypeLoader.findEntityId(ctx, entityType, entityName);
return Futures.immediateFuture(targetEntity);
return immediateFuture(targetEntity);
} catch (IllegalStateException e) {
return Futures.immediateFailedFuture(e);
}

50
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java

@ -1,50 +0,0 @@
/**
* Copyright © 2016-2025 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.rule.engine.util;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.HasCustomerId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UserId;
public class EntitiesCustomerIdAsyncLoader {
public static ListenableFuture<CustomerId> findEntityIdAsync(TbContext ctx, EntityId originator) {
switch (originator.getEntityType()) {
case CUSTOMER:
return Futures.immediateFuture((CustomerId) originator);
case USER:
return toCustomerIdAsync(ctx, ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originator));
case ASSET:
return toCustomerIdAsync(ctx, ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originator));
case DEVICE:
return toCustomerIdAsync(ctx, Futures.immediateFuture(ctx.getDeviceService().findDeviceById(ctx.getTenantId(), (DeviceId) originator)));
default:
return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originator.getEntityType()));
}
}
private static <T extends HasCustomerId> ListenableFuture<CustomerId> toCustomerIdAsync(TbContext ctx, ListenableFuture<T> future) {
return Futures.transform(future, in -> in != null ? in.getCustomerId() : null, ctx.getDbCallbackExecutor());
}
}

2
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/TestDbCallbackExecutor.java

@ -28,7 +28,7 @@ public class TestDbCallbackExecutor implements ListeningExecutor {
try {
return Futures.immediateFuture(task.call());
} catch (Exception e) {
throw new RuntimeException(e);
return Futures.immediateFailedFuture(e);
}
}

111
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java

@ -16,7 +16,7 @@
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.FluentFuture;
import lombok.RequiredArgsConstructor;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
@ -53,19 +53,19 @@ import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.user.UserService;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Optional;
import java.util.UUID;
import static com.google.common.util.concurrent.Futures.immediateFuture;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
@ -73,19 +73,20 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class TbGetCustomerAttributeNodeTest {
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID());
private static final TenantId TENANT_ID = new TenantId(UUID.randomUUID());
private static final CustomerId CUSTOMER_ID = new CustomerId(UUID.randomUUID());
private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
private final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID());
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.randomUUID());
private final CustomerId CUSTOMER_ID = new CustomerId(UUID.randomUUID());
private final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
@Mock
private TbContext ctxMock;
@Mock
@ -93,21 +94,24 @@ public class TbGetCustomerAttributeNodeTest {
@Mock
private TimeseriesService timeseriesServiceMock;
@Mock
private UserService userServiceMock;
@Mock
private AssetService assetServiceMock;
@Mock
private DeviceService deviceServiceMock;
private EntityService entityServiceMock;
private TbGetCustomerAttributeNode node;
private TbGetEntityDataNodeConfiguration config;
private TbNodeConfiguration nodeConfiguration;
private TbMsg msg;
@BeforeEach
public void setUp() {
public void setup() {
node = new TbGetCustomerAttributeNode();
config = new TbGetEntityDataNodeConfiguration().defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
lenient().when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
lenient().when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
lenient().when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock);
lenient().when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock);
lenient().when(ctxMock.getEntityService()).thenReturn(entityServiceMock);
}
@Test
@ -237,12 +241,9 @@ public class TbGetCustomerAttributeNodeTest {
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(ctxMock.getUserService()).thenReturn(userServiceMock);
doReturn(Futures.immediateFuture(null)).when(userServiceMock).findUserByIdAsync(eq(TENANT_ID), eq(userId));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, userId)).thenReturn(
FluentFuture.from(immediateFuture(Optional.empty()))
);
// WHEN
node.onMsg(ctxMock, msg);
@ -252,18 +253,13 @@ public class TbGetCustomerAttributeNodeTest {
var actualExceptionCaptor = ArgumentCaptor.forClass(Throwable.class);
verify(ctxMock, never()).tellSuccess(any());
verify(ctxMock, times(1))
.tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture());
verify(ctxMock).tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture());
var actualMessage = actualMessageCaptor.getValue();
var actualException = actualExceptionCaptor.getValue();
var expectedExceptionMessage = String.format(
"Failed to find customer for entity with id: %s and type: %s",
userId.getId(), userId.getEntityType().getNormalName());
assertEquals(msg, actualMessage);
assertEquals(expectedExceptionMessage, actualException.getMessage());
assertEquals("Originator not found", actualException.getMessage());
assertInstanceOf(NoSuchElementException.class, actualException);
}
@ -282,16 +278,12 @@ public class TbGetCustomerAttributeNodeTest {
);
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3");
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
doReturn(device).when(deviceServiceMock).findDeviceById(eq(TENANT_ID), eq(device.getId()));
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, device.getId())).thenReturn(
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID)))
);
when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock);
when(attributesServiceMock.find(eq(TENANT_ID), eq(CUSTOMER_ID), eq(AttributeScope.SERVER_SCOPE), argThat(new ListMatcher<>(expectedPatternProcessedKeysList))))
.thenReturn(Futures.immediateFuture(attributesList));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
.thenReturn(immediateFuture(attributesList));
// WHEN
node.onMsg(ctxMock, msg);
@ -299,7 +291,7 @@ public class TbGetCustomerAttributeNodeTest {
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgData = "{\"temp\":42," +
@ -329,16 +321,12 @@ public class TbGetCustomerAttributeNodeTest {
);
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3");
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(ctxMock.getUserService()).thenReturn(userServiceMock);
doReturn(Futures.immediateFuture(user)).when(userServiceMock).findUserByIdAsync(eq(TENANT_ID), eq(user.getId()));
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, user.getId())).thenReturn(
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID)))
);
when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock);
when(attributesServiceMock.find(eq(TENANT_ID), eq(CUSTOMER_ID), eq(AttributeScope.SERVER_SCOPE), argThat(new ListMatcher<>(expectedPatternProcessedKeysList))))
.thenReturn(Futures.immediateFuture(attributesList));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
.thenReturn(immediateFuture(attributesList));
// WHEN
node.onMsg(ctxMock, msg);
@ -346,7 +334,7 @@ public class TbGetCustomerAttributeNodeTest {
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgMetaData = new TbMsgMetaData(Map.of(
@ -375,21 +363,18 @@ public class TbGetCustomerAttributeNodeTest {
);
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3");
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock);
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(customer.getId()), argThat(new ListMatcher<>(expectedPatternProcessedKeysList))))
.thenReturn(Futures.immediateFuture(timeseriesList));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
.thenReturn(immediateFuture(timeseriesList));
// WHEN
node.onMsg(ctxMock, msg);
// THEN
verifyNoInteractions(entityServiceMock);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgData = "{\"temp\":42," +
@ -408,7 +393,7 @@ public class TbGetCustomerAttributeNodeTest {
public void givenFetchTelemetryToMetaData_whenOnMsg_thenShouldFetchTelemetryToMetaData() {
// GIVEN
var asset = new Asset(new AssetId(UUID.randomUUID()));
asset.setCustomerId(new CustomerId(UUID.randomUUID()));
asset.setCustomerId(CUSTOMER_ID);
prepareMsgAndConfig(TbMsgSource.METADATA, DataToFetch.LATEST_TELEMETRY, asset.getId());
@ -419,16 +404,12 @@ public class TbGetCustomerAttributeNodeTest {
);
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3");
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(ctxMock.getAssetService()).thenReturn(assetServiceMock);
doReturn(Futures.immediateFuture(asset)).when(assetServiceMock).findAssetByIdAsync(eq(TENANT_ID), eq(asset.getId()));
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock);
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(asset.getCustomerId()), argThat(new ListMatcher<>(expectedPatternProcessedKeysList))))
.thenReturn(Futures.immediateFuture(timeseriesList));
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, asset.getId())).thenReturn(
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID)))
);
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(CUSTOMER_ID), argThat(new ListMatcher<>(expectedPatternProcessedKeysList))))
.thenReturn(immediateFuture(timeseriesList));
// WHEN
node.onMsg(ctxMock, msg);
@ -436,7 +417,7 @@ public class TbGetCustomerAttributeNodeTest {
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgMetaData = new TbMsgMetaData(Map.of(

90
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java

@ -15,8 +15,10 @@
*/
package org.thingsboard.rule.engine.transform;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.Futures;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
@ -35,6 +37,7 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.data.RelationsQuery;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.alarm.Alarm;
@ -54,21 +57,24 @@ import org.thingsboard.server.common.data.relation.RelationsSearchParameters;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.relation.RelationService;
import java.util.Collections;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Optional;
import java.util.UUID;
import java.util.stream.Stream;
import static com.google.common.util.concurrent.Futures.immediateFuture;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.lenient;
import static org.thingsboard.rule.engine.transform.OriginatorSource.ALARM_ORIGINATOR;
import static org.thingsboard.rule.engine.transform.OriginatorSource.CUSTOMER;
import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY;
@ -93,16 +99,23 @@ public class TbChangeOriginatorNodeTest {
@Mock
private AssetService assetServiceMock;
@Mock
private DeviceService deviceServiceMock;
@Mock
private RelationService relationServiceMock;
@Mock
private RuleEngineAlarmService alarmServiceMock;
@Mock
private EntityService entityServiceMock;
@BeforeEach
public void before() throws TbNodeException {
public void setup() {
node = new TbChangeOriginatorNode();
config = new TbChangeOriginatorNodeConfiguration().defaultConfiguration();
lenient().when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
lenient().when(ctxMock.getDbCallbackExecutor()).thenReturn(dbExecutor);
lenient().when(ctxMock.getAssetService()).thenReturn(assetServiceMock);
lenient().when(ctxMock.getRelationService()).thenReturn(relationServiceMock);
lenient().when(ctxMock.getAlarmService()).thenReturn(alarmServiceMock);
lenient().when(ctxMock.getEntityService()).thenReturn(entityServiceMock);
}
@Test
@ -172,8 +185,13 @@ public class TbChangeOriginatorNodeTest {
}
@Test
public void givenOriginatorSourceIsCustomer_whenOnMsg_thenTellSuccess() throws TbNodeException {
Device device = new Device(DEVICE_ID);
@DisplayName("""
Given a device assigned to a customer and node configured to change originator to customer,
when processing the message,
then should change message originator from device to customer""")
public void givenDeviceAssignedToCustomer_whenProcessingMessage_thenChangesOriginatorToCustomer() throws TbNodeException {
// GIVEN
var device = new Device(DEVICE_ID);
device.setCustomerId(CUSTOMER_ID);
TbMsg msg = TbMsg.newMsg()
@ -186,16 +204,52 @@ public class TbChangeOriginatorNodeTest {
.originator(CUSTOMER_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getDeviceService()).willReturn(deviceServiceMock);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
given(deviceServiceMock.findDeviceById(any(TenantId.class), any(DeviceId.class))).willReturn(device);
given(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, device.getId())).willReturn(
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID)))
);
given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// WHEN
node.onMsg(ctxMock, msg);
// THEN
then(ctxMock).should().transformMsgOriginator(msg, CUSTOMER_ID);
ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellSuccess(actualMsg.capture());
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
}
@Test
@DisplayName("""
Given a customer as message originator and node configured to change originator to customer,
when processing the message,
then should keep the customer as originator""")
public void givenCustomerAsOriginator_whenProcessingMessage_thenKeepsCustomerAsOriginator() throws TbNodeException {
// GIVEN
var customer = new Customer(CUSTOMER_ID);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(customer.getId())
.metaData(TbMsgMetaData.EMPTY)
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
var expectedMsg = msg.transform()
.originator(CUSTOMER_ID)
.build();
given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// WHEN
node.onMsg(ctxMock, msg);
then(deviceServiceMock).should().findDeviceById(TENANT_ID, DEVICE_ID);
// THEN
then(ctxMock).should().transformMsgOriginator(msg, CUSTOMER_ID);
ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellSuccess(actualMsg.capture());
@ -216,8 +270,6 @@ public class TbChangeOriginatorNodeTest {
.originator(TENANT_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -240,9 +292,6 @@ public class TbChangeOriginatorNodeTest {
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getRelationService()).willReturn(relationServiceMock);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
given(relationServiceMock.findByQuery(any(TenantId.class), any(EntityRelationsQuery.class))).willReturn(Futures.immediateFuture(Collections.emptyList()));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -282,9 +331,6 @@ public class TbChangeOriginatorNodeTest {
.originator(DEVICE_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getAlarmService()).willReturn(alarmServiceMock);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
given(alarmServiceMock.findAlarmByIdAsync(any(TenantId.class), any(AlarmId.class))).willReturn(Futures.immediateFuture(alarm));
given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg);
@ -315,9 +361,6 @@ public class TbChangeOriginatorNodeTest {
.originator(ASSET_ID)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getAssetService()).willReturn(assetServiceMock);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
given(assetServiceMock.findAssetByTenantIdAndName(any(TenantId.class), any(String.class))).willReturn(new Asset(ASSET_ID));
given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg);
@ -355,9 +398,6 @@ public class TbChangeOriginatorNodeTest {
.data(TbMsg.EMPTY_JSON_OBJECT)
.build();
given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor);
given(ctxMock.getAssetService()).willReturn(assetServiceMock);
given(ctxMock.getTenantId()).willReturn(TENANT_ID);
given(assetServiceMock.findAssetByTenantIdAndName(any(TenantId.class), any(String.class))).willReturn(null);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));

154
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java

@ -1,154 +0,0 @@
/**
* Copyright © 2016-2025 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.rule.engine.util;
import com.google.common.util.concurrent.Futures;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.TestDbCallbackExecutor;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.user.UserService;
import java.util.EnumSet;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class EntitiesCustomerIdAsyncLoaderTest {
private static final EnumSet<EntityType> SUPPORTED_ENTITY_TYPES = EnumSet.of(
EntityType.CUSTOMER,
EntityType.USER,
EntityType.ASSET,
EntityType.DEVICE
);
private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
@Mock
private TbContext ctxMock;
@Mock
private UserService userServiceMock;
@Mock
private AssetService assetServiceMock;
@Mock
private DeviceService deviceServiceMock;
@Test
public void givenCustomerEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException {
// GIVEN
var customer = new Customer(new CustomerId(UUID.randomUUID()));
// WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, customer.getId()).get();
// THEN
assertEquals(customer.getId(), actualCustomerId);
}
@Test
public void givenUserEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException {
// GIVEN
var user = new User(new UserId(UUID.randomUUID()));
var expectedCustomerId = new CustomerId(UUID.randomUUID());
user.setCustomerId(expectedCustomerId);
when(ctxMock.getUserService()).thenReturn(userServiceMock);
doReturn(Futures.immediateFuture(user)).when(userServiceMock).findUserByIdAsync(any(), any());
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, user.getId()).get();
// THEN
assertEquals(expectedCustomerId, actualCustomerId);
}
@Test
public void givenAssetEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException {
// GIVEN
var asset = new Asset(new AssetId(UUID.randomUUID()));
var expectedCustomerId = new CustomerId(UUID.randomUUID());
asset.setCustomerId(expectedCustomerId);
when(ctxMock.getAssetService()).thenReturn(assetServiceMock);
doReturn(Futures.immediateFuture(asset)).when(assetServiceMock).findAssetByIdAsync(any(), any());
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, asset.getId()).get();
// THEN
assertEquals(expectedCustomerId, actualCustomerId);
}
@Test
public void givenDeviceEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException {
// GIVEN
var device = new Device(new DeviceId(UUID.randomUUID()));
var expectedCustomerId = new CustomerId(UUID.randomUUID());
device.setCustomerId(expectedCustomerId);
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
doReturn(device).when(deviceServiceMock).findDeviceById(any(), any());
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, device.getId()).get();
// THEN
assertEquals(expectedCustomerId, actualCustomerId);
}
@Test
public void givenUnsupportedEntityTypes_whenFindEntityIdAsync_thenException() {
for (var entityType : EntityType.values()) {
if (!SUPPORTED_ENTITY_TYPES.contains(entityType)) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, UUID.randomUUID());
var expectedExceptionMsg = "org.thingsboard.rule.engine.api.TbNodeException: Unexpected originator EntityType: " + entityType;
var exception = assertThrows(ExecutionException.class,
() -> EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, entityId).get());
assertInstanceOf(TbNodeException.class, exception.getCause());
assertEquals(expectedExceptionMsg, exception.getMessage());
}
}
}
}
Loading…
Cancel
Save