diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java index 59ae0c1953..0ab6cdf898 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationController.java @@ -33,8 +33,8 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.NotificationId; import org.thingsboard.server.common.data.id.NotificationRequestId; +import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.notification.Notification; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -148,8 +148,8 @@ public class NotificationController extends BaseController { preview.setProcessedTemplates(processedTemplates); Map recipientsCountByTarget = notificationRequest.getTargets().stream() - .collect(Collectors.toMap(UUIDBased::getId, targetId -> { - return notificationTargetService.countRecipientsForNotificationTarget(user.getTenantId(), targetId); + .collect(Collectors.toMap(id -> id, targetId -> { + return notificationTargetService.countRecipientsForNotificationTarget(user.getTenantId(), new NotificationTargetId(targetId)); })); preview.setRecipientsCountByTarget(recipientsCountByTarget); preview.setTotalRecipientsCount(recipientsCountByTarget.values().stream().mapToInt(Integer::intValue).sum()); diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java index 502bded58c..e8edac2e6c 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java @@ -120,7 +120,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple } } - notificationRequest.setStatus(NotificationRequestStatus.SENT); + notificationRequest.setStatus(NotificationRequestStatus.PROCESSING); NotificationRequest savedNotificationRequest = notificationRequestService.saveNotificationRequest(tenantId, notificationRequest); NotificationProcessingContext ctx = NotificationProcessingContext.builder() @@ -134,9 +134,9 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple Set deliveryMethods = ctx.getDeliveryMethods(); List> results = new ArrayList<>(); - for (NotificationTargetId targetId : notificationRequest.getTargets()) { + for (UUID targetId : notificationRequest.getTargets()) { DaoUtil.processBatches(pageLink -> { - return notificationTargetService.findRecipientsForNotificationTarget(tenantId, ctx.getCustomerId(), targetId, pageLink); + return notificationTargetService.findRecipientsForNotificationTarget(tenantId, ctx.getCustomerId(), new NotificationTargetId(targetId), pageLink); }, 200, recipientsBatch -> { for (NotificationDeliveryMethod deliveryMethod : deliveryMethods) { if (deliveryMethod.isIndependent()) continue; @@ -172,7 +172,8 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple Futures.whenAllComplete(results).run(() -> { NotificationRequestStats stats = ctx.getStats(); try { - notificationRequestService.updateNotificationRequestStats(tenantId, savedNotificationRequest.getId(), stats); + notificationRequestService.updateNotificationRequest(tenantId, savedNotificationRequest.getId(), + NotificationRequestStatus.SENT, stats); } catch (Exception e) { log.error("Failed to update stats for notification request {}", savedNotificationRequest.getId(), e); } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java index 588a333caf..9a32a5d458 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java @@ -30,6 +30,7 @@ import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.notification.info.AlarmOriginatedNotificationInfo; import org.thingsboard.server.common.data.notification.info.NotificationInfo; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -46,6 +47,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.executors.NotificationExecutorService; import java.util.List; +import java.util.stream.Collectors; @Service @TbCoreComponent @@ -131,7 +133,7 @@ public class DefaultNotificationRuleProcessingService implements NotificationRul NotificationInfo notificationInfo = constructNotificationInfo(alarm); NotificationRequest notificationRequest = NotificationRequest.builder() .tenantId(tenantId) - .targets(List.of(targetId)) + .targets(List.of(targetId).stream().map(UUIDBased::getId).collect(Collectors.toList())) .templateId(notificationRule.getTemplateId()) .additionalConfig(config) .info(notificationInfo) diff --git a/application/src/test/java/org/thingsboard/server/service/notification/AbstractNotificationApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/AbstractNotificationApiTest.java index 1ddcc9b8b9..90d81af64c 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/AbstractNotificationApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/AbstractNotificationApiTest.java @@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.NotificationTemplateId; +import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.notification.Notification; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; @@ -53,6 +54,7 @@ import java.net.URISyntaxException; import java.util.HashMap; import java.util.List; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; import static org.assertj.core.api.Assertions.assertThat; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -107,7 +109,7 @@ public abstract class AbstractNotificationApiTest extends AbstractControllerTest UserOriginatedNotificationInfo notificationInfo = new UserOriginatedNotificationInfo(); notificationInfo.setDescription("My description"); NotificationRequest notificationRequest = NotificationRequest.builder() - .targets(targets) + .targets(targets.stream().map(UUIDBased::getId).collect(Collectors.toList())) .templateId(notificationTemplateId) .info(notificationInfo) .additionalConfig(config) diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java index c8be17d3b1..820650d911 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java @@ -281,7 +281,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { @Test public void testNotificationUpdatesForALotOfUsers() throws Exception { - int usersCount = 200; // FIXME: sometimes if set e.g. to 150, up to 5 WS sessions don't receive update + int usersCount = 100; Map sessions = new HashMap<>(); List targets = new ArrayList<>(); @@ -331,7 +331,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { }); await().atMost(2, TimeUnit.SECONDS) - .until(() -> getStats(notificationRequest.getId()) != null); + .until(() -> findNotificationRequest(notificationRequest.getId()).isSent()); NotificationRequestStats stats = getStats(notificationRequest.getId()); assertThat(stats.getSent().get(NotificationDeliveryMethod.PUSH)).hasValue(usersCount); assertThat(stats.getSent().get(NotificationDeliveryMethod.EMAIL)).hasValue(usersCount); @@ -417,7 +417,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { NotificationRequest notificationRequest = new NotificationRequest(); - notificationRequest.setTargets(List.of(target1.getId(), target2.getId())); + notificationRequest.setTargets(List.of(target1.getUuidId(), target2.getUuidId())); notificationRequest.setTemplateId(notificationTemplate.getId()); notificationRequest.setAdditionalConfig(new NotificationRequestConfig()); @@ -471,7 +471,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { wsClient.waitForUpdate(); await().atMost(2, TimeUnit.SECONDS) - .until(() -> getStats(notificationRequest.getId()) != null); + .until(() -> findNotificationRequest(notificationRequest.getId()).isSent()); NotificationRequestStats stats = getStats(notificationRequest.getId()); assertThat(stats.getSent().get(NotificationDeliveryMethod.PUSH)).hasValue(1); @@ -555,7 +555,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { NotificationRequest successfulNotificationRequest = submitNotificationRequest(Collections.emptyList(), notificationTemplate.getId(), 0); await().atMost(2, TimeUnit.SECONDS) - .until(() -> getStats(successfulNotificationRequest.getId()) != null); + .until(() -> findNotificationRequest(successfulNotificationRequest.getId()).isSent()); verify(slackService).sendMessage(eq(tenantId), eq(slackToken), eq(conversationId), eq(config.getDefaultTextTemplate())); NotificationRequestStats stats = getStats(successfulNotificationRequest.getId()); assertThat(stats.getSent().get(NotificationDeliveryMethod.SLACK)).hasValue(1); @@ -564,7 +564,7 @@ public class NotificationApiTest extends AbstractNotificationApiTest { doThrow(new RuntimeException(errorMessage)).when(slackService).sendMessage(any(), any(), any(), any()); NotificationRequest failedNotificationRequest = submitNotificationRequest(Collections.emptyList(), notificationTemplate.getId(), 0); await().atMost(2, TimeUnit.SECONDS) - .until(() -> getStats(failedNotificationRequest.getId()) != null); + .until(() -> findNotificationRequest(failedNotificationRequest.getId()).isSent()); stats = getStats(failedNotificationRequest.getId()); assertThat(stats.getErrors().get(NotificationDeliveryMethod.SLACK).values()).containsExactly(errorMessage); } 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 db847a3e73..cbfbe8a408 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 @@ -44,6 +44,6 @@ public interface NotificationRequestService { PageData findScheduledNotificationRequests(PageLink pageLink); - void updateNotificationRequestStats(TenantId tenantId, NotificationRequestId notificationRequestId, NotificationRequestStats stats); + void updateNotificationRequest(TenantId tenantId, NotificationRequestId requestId, NotificationRequestStatus requestStatus, NotificationRequestStats stats); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java index e8e9e967e2..852adbdc60 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java @@ -27,16 +27,15 @@ import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.NotificationRuleId; -import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.NotificationTemplateId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.notification.info.NotificationInfo; import javax.validation.Valid; -import javax.validation.constraints.NotEmpty; import javax.validation.constraints.NotNull; import java.util.List; +import java.util.UUID; @Data @EqualsAndHashCode(callSuper = true) @@ -47,7 +46,7 @@ public class NotificationRequest extends BaseData impleme private TenantId tenantId; @NotNull - private List targets; + private List targets; @NotNull private NotificationTemplateId templateId; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStatus.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStatus.java index 7fedd8cdb1..d348fcfde4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStatus.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStatus.java @@ -16,6 +16,7 @@ package org.thingsboard.server.common.data.notification; public enum NotificationRequestStatus { + PROCESSING, SENT, SCHEDULED } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java index 2b8f1b32ad..9cf631a118 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java @@ -109,7 +109,7 @@ public class NotificationRequestEntity extends BaseSqlEntity new NotificationTargetId(UUID.fromString(uuid)))); + notificationRequest.setTargets(listFromString(targets, UUID::fromString)); notificationRequest.setTemplateId(getEntityId(templateId, NotificationTemplateId::new)); notificationRequest.setInfo(fromJson(info, NotificationInfo.class)); notificationRequest.setAdditionalConfig(fromJson(additionalConfig, NotificationRequestConfig.class)); 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 5530c4774e..90a7b5148a 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 @@ -79,8 +79,8 @@ public class DefaultNotificationRequestService implements NotificationRequestSer } @Override - public void updateNotificationRequestStats(TenantId tenantId, NotificationRequestId notificationRequestId, NotificationRequestStats stats) { - notificationRequestDao.updateStatsById(tenantId, notificationRequestId, stats); + public void updateNotificationRequest(TenantId tenantId, NotificationRequestId requestId, NotificationRequestStatus requestStatus, NotificationRequestStats stats) { + notificationRequestDao.updateById(tenantId, requestId, requestStatus, stats); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java index f415f2737a..ee1e6ab683 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java @@ -41,7 +41,7 @@ public interface NotificationRequestDao extends Dao { PageData findAllByStatus(NotificationRequestStatus status, PageLink pageLink); - void updateStatsById(TenantId tenantId, NotificationRequestId notificationRequestId, NotificationRequestStats stats); + void updateById(TenantId tenantId, NotificationRequestId requestId, NotificationRequestStatus requestStatus, NotificationRequestStats stats); boolean existsByStatusAndTargetId(TenantId tenantId, NotificationRequestStatus status, NotificationTargetId targetId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java index 0b869efd89..2b30d5d8fe 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java @@ -72,8 +72,8 @@ public class JpaNotificationRequestDao extends JpaAbstractDao { @NotEmpty - private List targets; + private List targets; @NotNull private NotificationTemplateId templateId; private NotificationRequestConfig additionalConfig;