diff --git a/application/src/main/java/org/thingsboard/server/controller/DashboardController.java b/application/src/main/java/org/thingsboard/server/controller/DashboardController.java index afc737ab88..7a1bb977ec 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DashboardController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DashboardController.java @@ -527,7 +527,7 @@ public class DashboardController extends BaseController { } Tenant tenant = tenantService.findTenantById(getTenantId()); JsonNode additionalInfo = tenant.getAdditionalInfo(); - if (additionalInfo == null || !(additionalInfo instanceof ObjectNode)) { + if (!(additionalInfo instanceof ObjectNode)) { additionalInfo = JacksonUtil.newObjectNode(); } if (homeDashboardInfo.getDashboardId() != null) { @@ -553,8 +553,7 @@ public class DashboardController extends BaseController { } return new HomeDashboardInfo(dashboardId, hideDashboardToolbar); } - } catch (Exception e) { - } + } catch (Exception ignored) {} return null; } @@ -570,8 +569,7 @@ public class DashboardController extends BaseController { } return new HomeDashboard(dashboard, hideDashboardToolbar); } - } catch (Exception e) { - } + } catch (Exception ignored) {} return null; } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index 69852e6c05..7106b8c94c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -239,7 +239,4 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { log.error("[{}] Can't push to edge updates, edgeNotificationMsg [{}]", tenantId, edgeNotificationMsg, throwable); callback.onFailure(throwable); } - } - - diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index aacacb711e..b077908bba 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -21,16 +21,12 @@ import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionalEventListener; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.cluster.TbClusterService; -import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.OtaPackageInfo; -import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.AlarmApiCallResult; import org.thingsboard.server.common.data.audit.ActionType; -import org.thingsboard.server.common.data.edge.Edge; -import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.relation.EntityRelation; @@ -79,7 +75,7 @@ public class EdgeEventSourcingListener { return; } try { - if (!isValidEdgeEventEntity(event.getEntity())) { + if (!isValidEdgeEventEntity(event)) { return; } log.trace("[{}] SaveEntityEvent called: {}", event.getTenantId(), event); @@ -150,24 +146,29 @@ public class EdgeEventSourcingListener { } } - private boolean isValidEdgeEventEntity(Object entity) { - if (entity instanceof OtaPackageInfo) { - OtaPackageInfo otaPackageInfo = (OtaPackageInfo) entity; - return otaPackageInfo.hasUrl() || otaPackageInfo.isHasData(); - } else if (entity instanceof RuleChain) { - RuleChain ruleChain = (RuleChain) entity; - return RuleChainType.EDGE.equals(ruleChain.getType()); - } else if (entity instanceof User) { - User user = (User) entity; - return !Authority.SYS_ADMIN.equals(user.getAuthority()); - } else if (entity instanceof AlarmApiCallResult) { - AlarmApiCallResult alarmApiCallResult = (AlarmApiCallResult) entity; - return alarmApiCallResult.isModified(); - } else if (entity instanceof Edge || - entity instanceof ApiUsageState || - entity instanceof TbResource || - entity instanceof EdgeEvent) { - return false; + private boolean isValidEdgeEventEntity(SaveEntityEvent event) { + switch (event.getEntityId().getEntityType()) { + case RULE_CHAIN: + RuleChain ruleChain = (RuleChain) event.getEntity(); + return RuleChainType.EDGE.equals(ruleChain.getType()); + case USER: + User user = (User) event.getEntity(); + return !Authority.SYS_ADMIN.equals(user.getAuthority()); + case OTA_PACKAGE: + OtaPackageInfo otaPackageInfo = (OtaPackageInfo) event.getEntity(); + return otaPackageInfo.hasUrl() || otaPackageInfo.isHasData(); + case ALARM: + if (event.getEntity() instanceof AlarmApiCallResult) { + AlarmApiCallResult alarmApiCallResult = (AlarmApiCallResult) event.getEntity(); + return alarmApiCallResult.isModified(); + } + break; + case TENANT: + return !event.getAdded(); + case API_USAGE_STATE: + case TB_RESOURCE: + case EDGE: + return false; } // Default: If the entity doesn't match any of the conditions, consider it as valid. return true; diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index 60a9f28c9a..250d88757a 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.entitiy; +import com.fasterxml.jackson.core.type.TypeReference; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; @@ -35,19 +36,25 @@ import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.entity.EntityStateSyncManager; import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; +import org.thingsboard.server.dao.tenant.TenantService; import javax.annotation.PostConstruct; +import java.util.Set; @Component @RequiredArgsConstructor @@ -55,6 +62,8 @@ import javax.annotation.PostConstruct; public class EntityStateSourcingListener { private final TbClusterService tbClusterService; + private final TenantService tenantService; + private final EntityStateSyncManager entityStateSyncManager; @PostConstruct public void init() { @@ -63,6 +72,9 @@ public class EntityStateSourcingListener { @TransactionalEventListener(fallbackExecution = true) public void handleEvent(SaveEntityEvent event) { + if (entityStateSyncManager.isSync()) { + return; + } log.trace("[{}] SaveEntityEvent called: {}", event.getTenantId(), event); TenantId tenantId = event.getTenantId(); EntityId entityId = event.getEntityId(); @@ -77,6 +89,12 @@ public class EntityStateSourcingListener { case NOTIFICATION_RULE: tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, lifecycleEvent); break; + case RULE_CHAIN: + RuleChain ruleChain = (RuleChain) event.getEntity(); + if (RuleChainType.CORE.equals(ruleChain.getType())) { + tbClusterService.broadcastEntityStateChangeEvent(ruleChain.getTenantId(), ruleChain.getId(), lifecycleEvent); + } + break; case TENANT: Tenant tenant = (Tenant) event.getEntity(); onTenantUpdate(tenant, lifecycleEvent); @@ -124,6 +142,17 @@ public class EntityStateSourcingListener { case NOTIFICATION_RULE: tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); break; + case RULE_CHAIN: + RuleChain ruleChain = (RuleChain) event.getEntity(); + Set referencingRuleChainIds = JacksonUtil.fromString(event.getBody(), new TypeReference<>() {}); + if (RuleChainType.CORE.equals(ruleChain.getType())) { + if (referencingRuleChainIds != null) { + referencingRuleChainIds.forEach(referencingRuleChainId -> + tbClusterService.broadcastEntityStateChangeEvent(tenantId, referencingRuleChainId, ComponentLifecycleEvent.UPDATED)); + } + tbClusterService.broadcastEntityStateChangeEvent(tenantId, ruleChain.getId(), ComponentLifecycleEvent.DELETED); + } + break; case TENANT: Tenant tenant = (Tenant) event.getEntity(); onTenantDeleted(tenant); @@ -154,11 +183,6 @@ public class EntityStateSourcingListener { } } - private void onDeviceProfileDelete(TenantId tenantId, EntityId entityId, DeviceProfile deviceProfile) { - tbClusterService.onDeviceProfileDelete(deviceProfile, null); - tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); - } - @TransactionalEventListener(fallbackExecution = true) public void handleEvent(ActionEntityEvent event) { log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event); @@ -178,6 +202,11 @@ public class EntityStateSourcingListener { tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), lifecycleEvent); } + private void onTenantDeleted(Tenant tenant) { + tbClusterService.onTenantDelete(tenant, null); + tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), ComponentLifecycleEvent.DELETED); + } + private void onTenantProfileUpdate(TenantProfile tenantProfile, ComponentLifecycleEvent lifecycleEvent) { tbClusterService.onTenantProfileChange(tenantProfile, null); tbClusterService.broadcastEntityStateChangeEvent(TenantId.SYS_TENANT_ID, tenantProfile.getId(), lifecycleEvent); @@ -198,6 +227,11 @@ public class EntityStateSourcingListener { return null; } + private void onDeviceProfileDelete(TenantId tenantId, EntityId entityId, DeviceProfile deviceProfile) { + tbClusterService.onDeviceProfileDelete(deviceProfile, null); + tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); + } + private void onDeviceUpdate(Object entity, Object oldEntity) { Device device = (Device) entity; Device oldDevice = null; @@ -207,11 +241,6 @@ public class EntityStateSourcingListener { tbClusterService.onDeviceUpdated(device, oldDevice); } - private void onTenantDeleted(Tenant tenant) { - tbClusterService.onTenantDelete(tenant, null); - tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), ComponentLifecycleEvent.DELETED); - } - private void handleEdgeEvent(TenantId tenantId, EntityId entityId, Object entity, ComponentLifecycleEvent lifecycleEvent) { if (entity instanceof Edge) { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, lifecycleEvent); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/tenant/DefaultTbTenantService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/tenant/DefaultTbTenantService.java index 8b1ffcfcf0..b9c8f29f82 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/tenant/DefaultTbTenantService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/tenant/DefaultTbTenantService.java @@ -16,10 +16,13 @@ package org.thingsboard.server.service.entitiy.tenant; import lombok.RequiredArgsConstructor; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.dao.entity.EntityStateSyncManager; +import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantProfileService; import org.thingsboard.server.dao.tenant.TenantService; @@ -43,12 +46,18 @@ public class DefaultTbTenantService extends AbstractTbEntityService implements T private final TbQueueService tbQueueService; private final TenantProfileService tenantProfileService; private final EntitiesVersionControlService versionControlService; + private final ApplicationEventPublisher eventPublisher; + private final EntityStateSyncManager entityStateSyncManager; @Override public Tenant save(Tenant tenant) throws Exception { boolean created = tenant.getId() == null; Tenant oldTenant = !created ? tenantService.findTenantById(tenant.getId()) : null; + if (created) { + entityStateSyncManager.getSync().set(true); + } + Tenant savedTenant = checkNotNull(tenantService.saveTenant(tenant)); if (created) { installScripts.createDefaultRuleChains(savedTenant.getId()); @@ -56,6 +65,11 @@ public class DefaultTbTenantService extends AbstractTbEntityService implements T } tenantProfileCache.evict(savedTenant.getId()); + if (created) { + entityStateSyncManager.getSync().remove(); + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(TenantId.SYS_TENANT_ID).entityId(savedTenant.getId()).entity(savedTenant).added(true).build()); + } + TenantProfile oldTenantProfile = oldTenant != null ? tenantProfileService.findTenantProfileById(TenantId.SYS_TENANT_ID, oldTenant.getTenantProfileId()) : null; TenantProfile newTenantProfile = tenantProfileService.findTenantProfileById(TenantId.SYS_TENANT_ID, savedTenant.getTenantProfileId()); tbQueueService.updateQueuesByTenants(Collections.singletonList(savedTenant.getTenantId()), newTenantProfile, oldTenantProfile); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/user/TbUserService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/user/TbUserService.java index 0764425116..08929d9ac9 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/user/TbUserService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/user/TbUserService.java @@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.id.TenantId; import javax.servlet.http.HttpServletRequest; public interface TbUserService { + User save(TenantId tenantId, CustomerId customerId, User tbUser, boolean sendActivationMail, HttpServletRequest request, User user) throws ThingsboardException; void delete(TenantId tenantId, CustomerId customerId, User user, User responsibleUser) throws ThingsboardException; diff --git a/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java b/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java index 7f0b70cb27..ec0d6fd268 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java +++ b/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java @@ -171,7 +171,7 @@ public class InstallScripts { return createRuleChainFromFile(tenantId, getDeviceProfileDefaultRuleChainTemplateFilePath(), ruleChainName); } - public RuleChain createRuleChainFromFile(TenantId tenantId, Path templateFilePath, String newRuleChainName) throws IOException { + public RuleChain createRuleChainFromFile(TenantId tenantId, Path templateFilePath, String newRuleChainName) { JsonNode ruleChainJson = JacksonUtil.toJsonNode(templateFilePath.toFile()); RuleChain ruleChain = JacksonUtil.treeToValue(ruleChainJson.get("ruleChain"), RuleChain.class); RuleChainMetaData ruleChainMetaData = JacksonUtil.treeToValue(ruleChainJson.get("metadata"), RuleChainMetaData.class); diff --git a/application/src/main/java/org/thingsboard/server/service/rule/DefaultTbRuleChainService.java b/application/src/main/java/org/thingsboard/server/service/rule/DefaultTbRuleChainService.java index b5d303a952..830c223d07 100644 --- a/application/src/main/java/org/thingsboard/server/service/rule/DefaultTbRuleChainService.java +++ b/application/src/main/java/org/thingsboard/server/service/rule/DefaultTbRuleChainService.java @@ -179,10 +179,6 @@ public class DefaultTbRuleChainService extends AbstractTbEntityService implement try { RuleChain savedRuleChain = checkNotNull(ruleChainService.saveRuleChain(ruleChain)); autoCommit(user, savedRuleChain.getId()); - if (RuleChainType.CORE.equals(savedRuleChain.getType())) { - tbClusterService.broadcastEntityStateChangeEvent(tenantId, savedRuleChain.getId(), - actionType.equals(ActionType.ADDED) ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - } logEntityActionService.logEntityAction(tenantId, savedRuleChain.getId(), savedRuleChain, null, actionType, user); return savedRuleChain; } catch (Exception e) { @@ -204,13 +200,6 @@ public class DefaultTbRuleChainService extends AbstractTbEntityService implement referencingRuleChainIds.remove(ruleChain.getId()); - if (RuleChainType.CORE.equals(ruleChain.getType())) { - referencingRuleChainIds.forEach(referencingRuleChainId -> - tbClusterService.broadcastEntityStateChangeEvent(tenantId, referencingRuleChainId, ComponentLifecycleEvent.UPDATED)); - - tbClusterService.broadcastEntityStateChangeEvent(tenantId, ruleChain.getId(), ComponentLifecycleEvent.DELETED); - } - logEntityActionService.logEntityAction(tenantId, ruleChainId, ruleChain, null, ActionType.DELETED, user, ruleChainId.toString()); } catch (Exception e) { logEntityActionService.logEntityAction(tenantId, emptyId(EntityType.RULE_CHAIN), ActionType.DELETED, @@ -224,7 +213,6 @@ public class DefaultTbRuleChainService extends AbstractTbEntityService implement try { RuleChain savedRuleChain = installScripts.createDefaultRuleChain(tenantId, request.getName()); autoCommit(user, savedRuleChain.getId()); - tbClusterService.broadcastEntityStateChangeEvent(tenantId, savedRuleChain.getId(), ComponentLifecycleEvent.CREATED); logEntityActionService.logEntityAction(tenantId, savedRuleChain.getId(), savedRuleChain, ActionType.ADDED, user); return savedRuleChain; } catch (Exception e) { @@ -243,18 +231,11 @@ public class DefaultTbRuleChainService extends AbstractTbEntityService implement RuleChain previousRootRuleChain = ruleChainService.getRootTenantRuleChain(tenantId); if (ruleChainService.setRootRuleChain(tenantId, ruleChainId)) { if (previousRootRuleChain != null) { - RuleChainId previousRootRuleChainId = previousRootRuleChain.getId(); - previousRootRuleChain = ruleChainService.findRuleChainById(tenantId, previousRootRuleChainId); - - tbClusterService.broadcastEntityStateChangeEvent(tenantId, previousRootRuleChainId, - ComponentLifecycleEvent.UPDATED); - logEntityActionService.logEntityAction(tenantId, previousRootRuleChainId, previousRootRuleChain, + previousRootRuleChain = ruleChainService.findRuleChainById(tenantId, previousRootRuleChain.getId()); + logEntityActionService.logEntityAction(tenantId, previousRootRuleChain.getId(), previousRootRuleChain, ActionType.UPDATED, user); } ruleChain = ruleChainService.findRuleChainById(tenantId, ruleChainId); - - tbClusterService.broadcastEntityStateChangeEvent(tenantId, ruleChainId, - ComponentLifecycleEvent.UPDATED); logEntityActionService.logEntityAction(tenantId, ruleChainId, ruleChain, ActionType.UPDATED, user); } return ruleChain; @@ -293,7 +274,6 @@ public class DefaultTbRuleChainService extends AbstractTbEntityService implement RuleChainMetaData savedRuleChainMetaData = checkNotNull(ruleChainService.loadRuleChainMetaData(tenantId, ruleChainMetaDataId)); if (RuleChainType.CORE.equals(ruleChain.getType())) { - tbClusterService.broadcastEntityStateChangeEvent(tenantId, ruleChainId, ComponentLifecycleEvent.UPDATED); updatedRuleChains.forEach(updatedRuleChain -> tbClusterService.broadcastEntityStateChangeEvent(tenantId, updatedRuleChain.getId(), ComponentLifecycleEvent.UPDATED)); } diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java index c2288a6f62..3732f13b08 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java +++ b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java @@ -24,7 +24,6 @@ import org.springframework.security.authentication.UsernamePasswordAuthenticatio import org.springframework.security.core.userdetails.UsernameNotFoundException; import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.StringUtils; @@ -88,7 +87,7 @@ public abstract class AbstractOAuth2ClientMapper { @Value("${edges.enabled}") @Getter private boolean edgesEnabled; - + private final Lock userCreationLock = new ReentrantLock(); protected SecurityUser getOrCreateSecurityUserFromOAuth2User(OAuth2User oauth2User, OAuth2Registration registration) { diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityStateSyncManager.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityStateSyncManager.java new file mode 100644 index 0000000000..69a4004498 --- /dev/null +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityStateSyncManager.java @@ -0,0 +1,23 @@ +/** + * Copyright © 2016-2023 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.server.dao.entity; + +public interface EntityStateSyncManager { + + ThreadLocal getSync(); + + boolean isSync(); +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java index 210810f68c..c785f88078 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java @@ -111,17 +111,13 @@ public class AssetProfileServiceImpl extends AbstractCachedEntityService sync = new ThreadLocal<>(); + + @Override + public boolean isSync() { + Boolean sync = this.sync.get(); + return sync != null && sync; + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java index 4335fa5f1c..ceecce8b1d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java @@ -101,17 +101,13 @@ public class EntityViewServiceImpl extends AbstractCachedEntityService { private final EntityId entityId; private final EdgeId edgeId; private final T entity; + private final String body; @Builder.Default private final long ts = System.currentTimeMillis(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 5e6b1e8009..c8024d31a2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -38,6 +38,7 @@ 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.page.PageData; +import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.plugin.ComponentClusteringMode; import org.thingsboard.server.common.data.relation.EntityRelation; @@ -73,6 +74,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @@ -135,6 +137,8 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC } else if (!previousRootRuleChain.getId().equals(ruleChain.getId())) { previousRootRuleChain.setRoot(false); ruleChainDao.save(tenantId, previousRootRuleChain); + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId) + .entityId(previousRootRuleChain.getId()).entity(previousRootRuleChain).added(false).build()); setRootAndSave(tenantId, ruleChain); return true; } @@ -145,6 +149,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC private void setRootAndSave(TenantId tenantId, RuleChain ruleChain) { ruleChain.setRoot(true); ruleChainDao.save(tenantId, ruleChain); + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entityId(ruleChain.getId()).entity(ruleChain).added(false).build()); } @Override @@ -290,11 +295,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC List nodeRelations = getRuleNodeRelations(tenantId, node.getId()); for (EntityRelation nodeRelation : nodeRelations) { String type = nodeRelation.getType(); - if (nodeRelation.getTo().getEntityType() == EntityType.RULE_NODE) { + if (EntityType.RULE_NODE.equals(nodeRelation.getTo().getEntityType())) { RuleNodeId toNodeId = new RuleNodeId(nodeRelation.getTo().getId()); int toIndex = ruleNodeIndexMap.get(toNodeId); ruleChainMetaData.addConnectionInfo(fromIndex, toIndex, type); - } else if (nodeRelation.getTo().getEntityType() == EntityType.RULE_CHAIN) { + } else if (EntityType.RULE_CHAIN.equals(nodeRelation.getTo().getEntityType())) { log.warn("[{}][{}] Unsupported node relation: {}", tenantId, ruleChainId, nodeRelation.getTo()); } } @@ -410,29 +415,24 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC public void deleteRuleChainById(TenantId tenantId, RuleChainId ruleChainId) { Validator.validateId(ruleChainId, "Incorrect rule chain id for delete request."); RuleChain ruleChain = ruleChainDao.findById(tenantId, ruleChainId.getId()); + + List referencingRuleNodes = getReferencingRuleChainNodes(tenantId, ruleChainId); + Set referencingRuleChainIds = referencingRuleNodes.stream().map(RuleNode::getRuleChainId).collect(Collectors.toSet()); + if (ruleChain != null) { if (ruleChain.isRoot()) { throw new DataValidationException("Deletion of Root Tenant Rule Chain is prohibited!"); } if (RuleChainType.EDGE.equals(ruleChain.getType())) { - PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); - PageData pageData; - do { - pageData = edgeService.findEdgesByTenantIdAndEntityId(tenantId, ruleChainId, pageLink); - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - for (Edge edge : pageData.getData()) { - if (edge.getRootRuleChainId() != null && edge.getRootRuleChainId().equals(ruleChainId)) { - throw new DataValidationException("Can't delete rule chain that is root for edge [" + edge.getName() + "]. Please assign another root rule chain first to the edge!"); - } - } - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } + for (Edge edge : new PageDataIterable<>(link -> edgeService.findEdgesByTenantIdAndEntityId(tenantId, ruleChainId, link), DEFAULT_PAGE_SIZE)) { + if (edge.getRootRuleChainId() != null && edge.getRootRuleChainId().equals(ruleChainId)) { + throw new DataValidationException("Can't delete rule chain that is root for edge [" + edge.getName() + "]. Please assign another root rule chain first to the edge!"); } - } while (pageData != null && pageData.hasNext()); + } } } - checkRuleNodesAndDelete(tenantId, ruleChainId); + checkRuleNodesAndDelete(tenantId, ruleChain, referencingRuleChainIds); + referencingRuleChainIds.remove(ruleChainId); } @Override @@ -731,11 +731,15 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC return ruleNodeDao.save(tenantId, ruleNode); } - private void checkRuleNodesAndDelete(TenantId tenantId, RuleChainId ruleChainId) { + private void checkRuleNodesAndDelete(TenantId tenantId, RuleChain ruleChain, Set referencingRuleChainIds) { try { entityCountService.publishCountEntityEvictEvent(tenantId, EntityType.RULE_CHAIN); - ruleChainDao.removeById(tenantId, ruleChainId.getId()); - eventPublisher.publishEvent(DeleteEntityEvent.builder().tenantId(tenantId).entityId(ruleChainId).build()); + ruleChainDao.removeById(tenantId, ruleChain.getUuidId()); + + if (referencingRuleChainIds != null) { + referencingRuleChainIds.remove(ruleChain.getId()); + } + eventPublisher.publishEvent(DeleteEntityEvent.builder().tenantId(tenantId).entityId(ruleChain.getId()).entity(ruleChain).body(JacksonUtil.toString(referencingRuleChainIds)).build()); } catch (Exception t) { ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("fk_default_rule_chain_device_profile")) { @@ -746,7 +750,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC throw t; } } - deleteRuleNodes(tenantId, ruleChainId); + deleteRuleNodes(tenantId, ruleChain.getId()); } private void deleteRuleNodes(TenantId tenantId, List ruleNodes) { @@ -815,7 +819,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Override protected void removeEntity(TenantId tenantId, RuleChain entity) { - checkRuleNodesAndDelete(tenantId, entity.getId()); + checkRuleNodesAndDelete(tenantId, entity, null); } };