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 80d2647078..98bd890fab 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 @@ -39,6 +39,7 @@ 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; @@ -47,7 +48,6 @@ 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.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; -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; @@ -61,7 +61,6 @@ import java.util.Set; public class EntityStateSourcingListener { private final TbClusterService tbClusterService; - private final EntityStateSyncManager entityStateSyncManager; @PostConstruct public void init() { @@ -70,9 +69,6 @@ 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(); @@ -138,13 +134,18 @@ public class EntityStateSourcingListener { case CUSTOMER: case EDGE: case NOTIFICATION_RULE: - case NOTIFICATION_REQUEST: tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); break; + case NOTIFICATION_REQUEST: + NotificationRequest request = (NotificationRequest) event.getEntity(); + if (request.isScheduled()) { + 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())) { + Set referencingRuleChainIds = JacksonUtil.fromString(event.getBody(), new TypeReference<>() {}); if (referencingRuleChainIds != null) { referencingRuleChainIds.forEach(referencingRuleChainId -> tbClusterService.broadcastEntityStateChangeEvent(tenantId, referencingRuleChainId, ComponentLifecycleEvent.UPDATED)); 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 5d8a477f5a..8201ded840 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 @@ -21,7 +21,6 @@ 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; @@ -47,17 +46,12 @@ public class DefaultTbTenantService extends AbstractTbEntityService implements T 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()); @@ -67,7 +61,6 @@ 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()); } 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 223fa39f3a..41392f242f 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 @@ -193,10 +193,10 @@ public class InstallScripts { if (!StringUtils.isEmpty(newRuleChainName)) { ruleChain.setName(newRuleChainName); } - ruleChain = ruleChainService.saveRuleChain(ruleChain); + ruleChain = ruleChainService.saveRuleChain(ruleChain, false); ruleChainMetaData.setRuleChainId(ruleChain.getId()); - ruleChainService.saveRuleChainMetaData(TenantId.SYS_TENANT_ID, ruleChainMetaData, Function.identity()); + ruleChainService.saveRuleChainMetaData(TenantId.SYS_TENANT_ID, ruleChainMetaData, Function.identity(), false); return ruleChain; } 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 49d33122ef..deee755cd3 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 @@ -42,7 +42,6 @@ import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.dashboard.DashboardService; -import org.thingsboard.server.dao.entity.EntityStateSyncManager; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.oauth2.OAuth2User; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; @@ -90,9 +89,6 @@ public abstract class AbstractOAuth2ClientMapper { @Autowired private ApplicationEventPublisher eventPublisher; - @Autowired - private EntityStateSyncManager entityStateSyncManager; - @Value("${edges.enabled}") @Getter private boolean edgesEnabled; @@ -179,8 +175,6 @@ public abstract class AbstractOAuth2ClientMapper { List tenants = tenantService.findTenants(new PageLink(1, 0, tenantName)).getData(); Tenant tenant; if (tenants == null || tenants.isEmpty()) { - entityStateSyncManager.getSync().set(true); - tenant = new Tenant(); tenant.setTitle(tenantName); tenant = tenantService.saveTenant(tenant); @@ -188,7 +182,6 @@ public abstract class AbstractOAuth2ClientMapper { installScripts.createDefaultEdgeRuleChains(tenant.getId()); tenantProfileCache.evict(tenant.getId()); - entityStateSyncManager.getSync().remove(); eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(TenantId.SYS_TENANT_ID).entityId(tenant.getId()).entity(tenant).added(true).build()); } else { tenant = tenants.get(0); 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 deleted file mode 100644 index 69a4004498..0000000000 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityStateSyncManager.java +++ /dev/null @@ -1,23 +0,0 @@ -/** - * 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/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index 532da4ac00..85f341b95e 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -44,10 +44,14 @@ public interface RuleChainService extends EntityDaoService { RuleChain saveRuleChain(RuleChain ruleChain); + RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent); + boolean setRootRuleChain(TenantId tenantId, RuleChainId ruleChainId); RuleChainUpdateResult saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData, Function ruleNodeUpdater); + RuleChainUpdateResult saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData, Function ruleNodeUpdater, boolean publishSaveEvent); + RuleChainMetaData loadRuleChainMetaData(TenantId tenantId, RuleChainId ruleChainId); RuleChain findRuleChainById(TenantId tenantId, RuleChainId ruleChainId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 59b0da9dae..68b77ce149 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -81,7 +81,6 @@ import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.tenant.TenantService; import java.util.ArrayList; -import java.util.Comparator; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -465,7 +464,7 @@ public class DeviceServiceImpl extends AbstractCachedEntityService customerDeviceUnasigner = new PaginatedRemover<>() { + private final PaginatedRemover customerDevicesRemover = new PaginatedRemover<>() { @Override protected PageData findEntities(TenantId tenantId, CustomerId id, PageLink pageLink) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/DefaultEntityStateSyncManager.java b/dao/src/main/java/org/thingsboard/server/dao/entity/DefaultEntityStateSyncManager.java deleted file mode 100644 index 5004b29995..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/DefaultEntityStateSyncManager.java +++ /dev/null @@ -1,34 +0,0 @@ -/** - * 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; - -import lombok.Getter; -import lombok.extern.slf4j.Slf4j; -import org.springframework.stereotype.Component; - -@Component -@Slf4j -public class DefaultEntityStateSyncManager implements EntityStateSyncManager { - - @Getter - private final ThreadLocal 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/notification/DefaultNotificationRequestService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java index 9f63cc71cf..740849826b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java @@ -90,9 +90,7 @@ public class DefaultNotificationRequestService implements NotificationRequestSer public void deleteNotificationRequest(TenantId tenantId, NotificationRequest request) { notificationRequestDao.removeById(tenantId, request.getUuidId()); notificationDao.deleteByRequestId(tenantId, request.getId()); - if (request.isScheduled()) { - eventPublisher.publishEvent(DeleteEntityEvent.builder().tenantId(tenantId).entityId(request.getId()).build()); - } + eventPublisher.publishEvent(DeleteEntityEvent.builder().tenantId(tenantId).entity(request).entityId(request.getId()).build()); } @Override 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 13ee23e9d5..d07102b80b 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 @@ -111,14 +111,22 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Override @Transactional public RuleChain saveRuleChain(RuleChain ruleChain) { + return saveRuleChain(ruleChain, true); + } + + @Override + @Transactional + public RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent) { ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); try { RuleChain savedRuleChain = ruleChainDao.save(ruleChain.getTenantId(), ruleChain); if (ruleChain.getId() == null) { entityCountService.publishCountEntityEvictEvent(ruleChain.getTenantId(), EntityType.RULE_CHAIN); } - eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(savedRuleChain.getTenantId()) - .entity(savedRuleChain).entityId(savedRuleChain.getId()).added(ruleChain.getId() == null).build()); + if (publishSaveEvent) { + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(savedRuleChain.getTenantId()) + .entity(savedRuleChain).entityId(savedRuleChain.getId()).added(ruleChain.getId() == null).build()); + } return savedRuleChain; } catch (Exception e) { checkConstraintViolation(e, "rule_chain_external_id_unq_key", "Rule Chain with such external id already exists!"); @@ -155,6 +163,12 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Override public RuleChainUpdateResult saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData, Function ruleNodeUpdater) { + return saveRuleChainMetaData(tenantId, ruleChainMetaData, ruleNodeUpdater, true); + } + + + @Override + public RuleChainUpdateResult saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData, Function ruleNodeUpdater, boolean publishSaveEvent) { Validator.validateId(ruleChainMetaData.getRuleChainId(), "Incorrect rule chain id."); RuleChain ruleChain = findRuleChainById(tenantId, ruleChainMetaData.getRuleChainId()); if (ruleChain == null) { @@ -268,8 +282,9 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC if (!relations.isEmpty()) { relationService.saveRelations(tenantId, relations); } - eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entity(ruleChain).entityId(ruleChain.getId()).build()); - + if (publishSaveEvent) { + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entity(ruleChain).entityId(ruleChain.getId()).build()); + } return RuleChainUpdateResult.successful(updatedRuleNodes); }