diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/DefaultTbAlarmService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/DefaultTbAlarmService.java index f1f8ba038c..48e496c05b 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/DefaultTbAlarmService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/DefaultTbAlarmService.java @@ -39,6 +39,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; +import org.thingsboard.server.dao.housekeeper.HouseKeeperService; import org.thingsboard.server.service.entitiy.AbstractTbEntityService; import java.util.List; @@ -51,6 +52,9 @@ public class DefaultTbAlarmService extends AbstractTbEntityService implements Tb @Autowired protected TbAlarmCommentService alarmCommentService; + @Autowired + private HouseKeeperService housekeeper; + @Override public Alarm save(Alarm alarm, User user) throws ThingsboardException { ActionType actionType = alarm.getId() == null ? ActionType.ADDED : ActionType.UPDATED; @@ -220,16 +224,10 @@ public class DefaultTbAlarmService extends AbstractTbEntityService implements Tb } @Override - public Boolean delete(Alarm alarm, User user) { - TenantId tenantId = alarm.getTenantId(); - notificationEntityService.logEntityAction(tenantId, alarm.getOriginator(), alarm, alarm.getCustomerId(), - ActionType.DELETED, user); - return alarmSubscriptionService.deleteAlarm(tenantId, alarm.getId()); - } - - private void unassignDeletedUserAlarms(User user) { + public List unassignDeletedUserAlarms(User user) { List alarmIds = alarmService.findAlarmIdsByAssigneeId(user.getId()); for (AlarmId alarmId : alarmIds) { + log.trace("[{}] Unassigning alarm {} userId {}", user.getTenantId().getId(), alarmId.getId(), user.getId().getId()); AlarmApiCallResult result = alarmSubscriptionService.unassignAlarm(user.getTenantId(), alarmId, System.currentTimeMillis()); Alarm alarm = result.getAlarm(); if (!result.isSuccessful()) { @@ -254,21 +252,26 @@ public class DefaultTbAlarmService extends AbstractTbEntityService implements Tb alarm.getCustomerId(), ActionType.ALARM_UNASSIGNED, null); } } + return alarmIds; } @TransactionalEventListener(fallbackExecution = true) public void handleEvent(DeleteEntityEvent event) { - try { - log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); - EntityId entityId = event.getEntityId(); - if (EntityType.USER.equals(entityId.getEntityType())) { - unassignDeletedUserAlarms((User) event.getEntity()); - } - } catch (Exception e) { - log.error("[{}] failed to process DeleteEntityEvent: {}", event.getTenantId(), event); + log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); + EntityId entityId = event.getEntityId(); + if (EntityType.USER.equals(entityId.getEntityType())) { + housekeeper.unassignDeletedUserAlarms((User) event.getEntity()); } } + @Override + public Boolean delete(Alarm alarm, User user) { + TenantId tenantId = alarm.getTenantId(); + notificationEntityService.logEntityAction(tenantId, alarm.getOriginator(), alarm, alarm.getCustomerId(), + ActionType.DELETED, user); + return alarmSubscriptionService.deleteAlarm(tenantId, alarm.getId()); + } + private static long getOrDefault(long ts) { return ts > 0 ? ts : System.currentTimeMillis(); } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/TbAlarmService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/TbAlarmService.java index a2ae9c8cc7..c531c20c53 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/TbAlarmService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/alarm/TbAlarmService.java @@ -19,8 +19,11 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.UserId; +import java.util.List; + public interface TbAlarmService { Alarm save(Alarm entity, User user) throws ThingsboardException; @@ -37,5 +40,7 @@ public interface TbAlarmService { AlarmInfo unassign(Alarm alarm, long unassignTs, User user) throws ThingsboardException; + List unassignDeletedUserAlarms(User user); + Boolean delete(Alarm alarm, User user); } diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/InMemoryHouseKeeperServiceService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/InMemoryHouseKeeperServiceService.java new file mode 100644 index 0000000000..2e62a68333 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/InMemoryHouseKeeperServiceService.java @@ -0,0 +1,73 @@ +/** + * 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.service.housekeeper; + +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningExecutorService; +import com.google.common.util.concurrent.MoreExecutors; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Lazy; +import org.springframework.stereotype.Component; +import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.dao.housekeeper.HouseKeeperService; +import org.thingsboard.server.service.entitiy.alarm.TbAlarmService; + +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.List; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicInteger; + +@Component +@RequiredArgsConstructor +@Slf4j +public class InMemoryHouseKeeperServiceService implements HouseKeeperService { + + @Lazy + final TbAlarmService alarmService; + + ListeningExecutorService executor; + + AtomicInteger queueSize = new AtomicInteger(); + + @PostConstruct + public void init() { + executor = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper"))); + } + + @PreDestroy + public void destroy() { + if (executor != null) { + executor.shutdown(); + } + } + + @Override + public ListenableFuture> unassignDeletedUserAlarms(User user) { + log.debug("[{}][{}] unassignDeletedUserAlarms submitting, pending queue size: {} ", user.getTenantId().getId(), user.getId().getId(), queueSize.get()); + queueSize.incrementAndGet(); + ListenableFuture> future = executor.submit(() -> alarmService.unassignDeletedUserAlarms(user)); + future.addListener(() -> { + queueSize.decrementAndGet(); + log.debug("[{}][{}] unassignDeletedUserAlarms finished, pending queue size: {} ", user.getTenantId().getId(), user.getId().getId(), queueSize.get()); + }, MoreExecutors.directExecutor()); + return future; + } + +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/housekeeper/HouseKeeperService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/housekeeper/HouseKeeperService.java new file mode 100644 index 0000000000..3ffd309ca5 --- /dev/null +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/housekeeper/HouseKeeperService.java @@ -0,0 +1,27 @@ +/** + * 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.housekeeper; + +import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.id.AlarmId; + +import java.util.List; + +public interface HouseKeeperService { + ListenableFuture> unassignDeletedUserAlarms(User user); + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java b/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java index 63b13bca66..9688f7b7e1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java +++ b/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java @@ -33,6 +33,7 @@ import java.util.Optional; import java.util.UUID; import java.util.function.Consumer; import java.util.function.Function; +import java.util.stream.Collectors; public abstract class DaoUtil { @@ -109,6 +110,10 @@ public abstract class DaoUtil { return ids; } + public static List fromUUIDs(List uuids, Function mapper) { + return uuids.stream().map(mapper).collect(Collectors.toList()); + } + public static I toEntityId(UUID uuid, Function creator) { if (uuid != null) { return creator.apply(uuid); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java index 045297f4ba..ed25595442 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java @@ -289,7 +289,7 @@ public class JpaAlarmDao extends JpaAbstractDao implements A @Override public List findAlarmIdsByAssigneeId(UUID key) { List assignedAlarmIds = alarmRepository.findAlarmIdsByAssigneeId(key); - return assignedAlarmIds.stream().map(AlarmId::new).collect(Collectors.toList()); + return DaoUtil.fromUUIDs(assignedAlarmIds, AlarmId::new); } @Override