From 92b09b85144a3b756901a7ba369fb4068576c37f Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 5 Oct 2023 14:09:03 +0300 Subject: [PATCH] Minor refactoring: still not working as expected --- .../NotificationRuleController.java | 8 +---- .../DefaultTbApiUsageStateService.java | 6 ++-- .../edge/EdgeEventSourcingListener.java | 4 +-- .../entitiy/EntityStateSourcingListener.java | 36 ++++++++----------- .../device/DefaultTbDeviceService.java | 5 +-- .../entitiy/user/DefaultUserService.java | 4 +-- .../impl/NotificationRuleImportService.java | 6 +--- .../DefaultTbAlarmCommentServiceTest.java | 6 ++-- .../server/dao/device/DeviceService.java | 3 +- .../NotificationRequestService.java | 2 +- .../server/dao/device/DeviceServiceImpl.java | 7 +++- .../server/dao/edge/BaseEdgeEventService.java | 27 +++++++++++--- .../usagerecord/ApiUsageStateServiceImpl.java | 2 +- 13 files changed, 58 insertions(+), 58 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java index 2afb67a8bd..822b6b44ab 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationRuleController.java @@ -37,7 +37,6 @@ import org.thingsboard.server.common.data.notification.rule.NotificationRuleInfo import org.thingsboard.server.common.data.notification.rule.trigger.config.NotificationRuleTriggerType; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.dao.notification.NotificationRuleService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.model.SecurityUser; @@ -125,11 +124,7 @@ public class NotificationRuleController extends BaseController { throw new IllegalArgumentException("Trigger type " + triggerType + " is not available"); } - boolean created = notificationRule.getId() == null; - notificationRule = doSaveAndLog(EntityType.NOTIFICATION_RULE, notificationRule, notificationRuleService::saveNotificationRule); - tbClusterService.broadcastEntityStateChangeEvent(user.getTenantId(), notificationRule.getId(), created ? - ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - return notificationRule; + return doSaveAndLog(EntityType.NOTIFICATION_RULE, notificationRule, notificationRuleService::saveNotificationRule); } @ApiOperation(value = "Get notification rule by id (getNotificationRuleById)", @@ -177,7 +172,6 @@ public class NotificationRuleController extends BaseController { NotificationRuleId notificationRuleId = new NotificationRuleId(id); NotificationRule notificationRule = checkEntityId(notificationRuleId, notificationRuleService::findNotificationRuleById, Operation.DELETE); doDeleteAndLog(EntityType.NOTIFICATION_RULE, notificationRule, notificationRuleService::deleteNotificationRuleById); - tbClusterService.broadcastEntityStateChangeEvent(user.getTenantId(), notificationRuleId, ComponentLifecycleEvent.DELETED); } } diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 2b1f7a0f6b..3f368f496d 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -25,7 +25,6 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.MailService; -import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.ApiFeature; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.ApiUsageRecordState; @@ -45,10 +44,11 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.notification.rule.trigger.ApiUsageLimitTrigger; import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfiguration; import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; -import org.thingsboard.server.common.data.notification.rule.trigger.ApiUsageLimitTrigger; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -61,7 +61,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceM import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.service.apiusage.BaseApiUsageState.StatsCalculationResult; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.mail.MailExecutorService; @@ -101,7 +100,6 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService public void onFailure(Throwable t) { } }; - private final TbClusterService clusterService; private final PartitionService partitionService; private final TenantService tenantService; private final TimeseriesService tsService; 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 57dba7e86e..aacacb711e 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 @@ -98,9 +98,7 @@ public class EdgeEventSourcingListener { } EntityType entityType = event.getEntityId().getEntityType(); try { - if (EntityType.EDGE.equals(entityType) || - EntityType.TENANT.equals(entityType) || - EntityType.TB_RESOURCE.equals(entityType)) { + if (EntityType.EDGE.equals(entityType) || EntityType.TENANT.equals(entityType)) { return; } log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); 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 6ed5c01f12..60a9f28c9a 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 @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; 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.security.DeviceCredentials; import org.thingsboard.server.common.msg.TbMsg; @@ -73,15 +74,16 @@ public class EntityStateSourcingListener { case ASSET: case ASSET_PROFILE: case ENTITY_VIEW: + case NOTIFICATION_RULE: tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, lifecycleEvent); break; case TENANT: Tenant tenant = (Tenant) event.getEntity(); - onTenantUpdate(tenant, isCreated); + onTenantUpdate(tenant, lifecycleEvent); break; case TENANT_PROFILE: TenantProfile tenantProfile = (TenantProfile) event.getEntity(); - onTenantProfileUpdate(tenantProfile, isCreated); + onTenantProfileUpdate(tenantProfile, lifecycleEvent); break; case DEVICE: onDeviceUpdate(event.getEntity(), event.getOldEntity()); @@ -118,6 +120,8 @@ public class EntityStateSourcingListener { case ASSET_PROFILE: case ENTITY_VIEW: case CUSTOMER: + case EDGE: + case NOTIFICATION_RULE: tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); break; case TENANT: @@ -136,13 +140,15 @@ public class EntityStateSourcingListener { DeviceProfile deviceProfile = (DeviceProfile) event.getEntity(); onDeviceProfileDelete(event.getTenantId(), event.getEntityId(), deviceProfile); break; - case EDGE: - tbClusterService.broadcastEntityStateChangeEvent(event.getTenantId(), event.getEntityId(), ComponentLifecycleEvent.DELETED); - break; case TB_RESOURCE: TbResource tbResource = (TbResource) event.getEntity(); tbClusterService.onResourceDeleted(tbResource, null); break; + case NOTIFICATION_REQUEST: + NotificationRequest request = (NotificationRequest) event.getEntity(); + if (request.isScheduled()) { + tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); + } default: break; } @@ -167,26 +173,14 @@ public class EntityStateSourcingListener { } } - private void onTenantUpdate(Tenant tenant, boolean isCreated) { + private void onTenantUpdate(Tenant tenant, ComponentLifecycleEvent lifecycleEvent) { tbClusterService.onTenantChange(tenant, null); - tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), isCreated ? - ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); + tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), lifecycleEvent); } - private void onTenantProfileUpdate(TenantProfile tenantProfile, boolean isCreated) { + private void onTenantProfileUpdate(TenantProfile tenantProfile, ComponentLifecycleEvent lifecycleEvent) { tbClusterService.onTenantProfileChange(tenantProfile, null); - tbClusterService.broadcastEntityStateChangeEvent(TenantId.SYS_TENANT_ID, tenantProfile.getId(), - isCreated ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - } - - private boolean isCommonEntityStateUpdated(EntityId entityId) { - switch (entityId.getEntityType()) { - case ASSET: - case ASSET_PROFILE: - case ENTITY_VIEW: - return true; - } - return false; + tbClusterService.broadcastEntityStateChangeEvent(TenantId.SYS_TENANT_ID, tenantProfile.getId(), lifecycleEvent); } private void onDeviceProfileUpdate(DeviceProfile deviceProfile, Object oldEntity, boolean isCreated) { diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/device/DefaultTbDeviceService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/device/DefaultTbDeviceService.java index ab5b5088fc..b6785d41d6 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/device/DefaultTbDeviceService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/device/DefaultTbDeviceService.java @@ -40,7 +40,6 @@ import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.claim.ClaimResponse; import org.thingsboard.server.dao.device.claim.ClaimResult; import org.thingsboard.server.dao.device.claim.ReclaimResult; -import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.entitiy.AbstractTbEntityService; @@ -53,7 +52,6 @@ public class DefaultTbDeviceService extends AbstractTbEntityService implements T private final DeviceService deviceService; private final DeviceCredentialsService deviceCredentialsService; private final ClaimDevicesService claimDevicesService; - private final TenantService tenantService; @Override public Device save(Device device, String accessToken, User user) throws Exception { @@ -230,8 +228,7 @@ public class DefaultTbDeviceService extends AbstractTbEntityService implements T TenantId newTenantId = newTenant.getId(); DeviceId deviceId = device.getId(); try { - Tenant tenant = tenantService.findTenantById(tenantId); - Device assignedDevice = deviceService.assignDeviceToTenant(tenant, newTenantId, device); + Device assignedDevice = deviceService.assignDeviceToTenant(newTenantId, device); logEntityActionService.logEntityAction(tenantId, deviceId, assignedDevice, assignedDevice.getCustomerId(), actionType, user, newTenantId.toString(), newTenant.getName()); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/user/DefaultUserService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/user/DefaultUserService.java index ab4354a709..7ef110cb29 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/user/DefaultUserService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/user/DefaultUserService.java @@ -81,10 +81,10 @@ public class DefaultUserService extends AbstractTbEntityService implements TbUse try { userService.deleteUser(tenantId, user); - logEntityActionService.logEntityAction(tenantId, userId, user, customerId, actionType, user, customerId.toString()); + logEntityActionService.logEntityAction(tenantId, userId, user, customerId, actionType, responsibleUser, customerId.toString()); } catch (Exception e) { logEntityActionService.logEntityAction(tenantId, emptyId(EntityType.USER), - actionType, user, e, userId.toString()); + actionType, responsibleUser, e, userId.toString()); throw e; } } diff --git a/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/impl/NotificationRuleImportService.java b/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/impl/NotificationRuleImportService.java index 9ace036ee6..0822640c4e 100644 --- a/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/impl/NotificationRuleImportService.java +++ b/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/impl/NotificationRuleImportService.java @@ -20,7 +20,6 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.audit.ActionType; -import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.NotificationRuleId; @@ -36,7 +35,6 @@ import org.thingsboard.server.common.data.notification.rule.trigger.config.Devic import org.thingsboard.server.common.data.notification.rule.trigger.config.NotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.config.NotificationRuleTriggerType; import org.thingsboard.server.common.data.notification.rule.trigger.config.RuleEngineComponentLifecycleEventNotificationRuleTriggerConfig; -import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.sync.ie.EntityExportData; import org.thingsboard.server.dao.notification.NotificationRuleService; import org.thingsboard.server.dao.service.ConstraintValidator; @@ -130,11 +128,9 @@ public class NotificationRuleImportService extends BaseEntityImportService> findDeviceTypesByTenantId(TenantId tenantId); - Device assignDeviceToTenant(Tenant oldTenant, TenantId tenantId, Device device); + Device assignDeviceToTenant(TenantId tenantId, Device device); PageData findDevicesIdsByDeviceProfileTransportType(DeviceTransportType transportType, PageLink pageLink); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java index af214dbf3d..22e7cfa530 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java @@ -45,7 +45,7 @@ public interface NotificationRequestService { List findNotificationRequestsByRuleIdAndOriginatorEntityId(TenantId tenantId, NotificationRuleId ruleId, EntityId originatorEntityId); - void deleteNotificationRequest(TenantId tenantId, NotificationRequestId requestId); + void deleteNotificationRequest(TenantId tenantId, NotificationRequest request); PageData findScheduledNotificationRequests(PageLink pageLink); 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 17975370ea..17da713dae 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 @@ -78,6 +78,7 @@ import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; +import org.thingsboard.server.dao.tenant.TenantService; import java.util.ArrayList; import java.util.Comparator; @@ -114,6 +115,9 @@ public class DeviceServiceImpl extends AbstractCachedEntityService deviceValidator; @@ -497,9 +501,10 @@ public class DeviceServiceImpl extends AbstractCachedEntityService entityViews = entityViewService.findEntityViewsByTenantIdAndEntityId(oldTenantId, device.getId()); if (!CollectionUtils.isEmpty(entityViews)) { throw new DataValidationException("Can't assign device that has entity views to another tenant!"); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java index fbb03ebb54..46c5671002 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java @@ -18,12 +18,12 @@ package org.thingsboard.server.dao.edge; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; -import lombok.AllArgsConstructor; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.jetbrains.annotations.NotNull; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; +import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; @@ -32,9 +32,14 @@ import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.service.DataValidator; +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + @Service @Slf4j -@AllArgsConstructor +@RequiredArgsConstructor public class BaseEdgeEventService implements EdgeEventService { private final EdgeEventDao edgeEventDao; @@ -43,6 +48,20 @@ public class BaseEdgeEventService implements EdgeEventService { private final ApplicationEventPublisher eventPublisher; + private ExecutorService executor; + + @PostConstruct + public void initExecutor() { + executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event")); + } + + @PreDestroy + public void shutdownExecutor() { + if (executor != null) { + executor.shutdown(); + } + } + @Override public ListenableFuture saveAsync(EdgeEvent edgeEvent) { edgeEventValidator.validate(edgeEvent, EdgeEvent::getTenantId); @@ -58,7 +77,7 @@ public class BaseEdgeEventService implements EdgeEventService { @Override public void onFailure(@NotNull Throwable throwable) {} - }, MoreExecutors.directExecutor()); + }, executor); return saveFuture; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java index 70d6e2bec1..90f00d0ad6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java @@ -146,7 +146,7 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A validateId(apiUsageState.getId(), "Can't save new usage state. Only update is allowed!"); ApiUsageState savedState = apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState); eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(apiUsageState.getTenantId()).entityId(savedState.getId()) - .entity(savedState).added(false).build()); + .entity(savedState).build()); return savedState; }