Browse Source

'PROCESSING' notification request status; targets to UUID array

pull/7980/head
ViacheslavKlimov 4 years ago
parent
commit
e4c2d41561
  1. 6
      application/src/main/java/org/thingsboard/server/controller/NotificationController.java
  2. 9
      application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java
  3. 4
      application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationRuleProcessingService.java
  4. 4
      application/src/test/java/org/thingsboard/server/service/notification/AbstractNotificationApiTest.java
  5. 12
      application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java
  6. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java
  7. 5
      common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequest.java
  8. 1
      common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStatus.java
  9. 2
      dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java
  10. 4
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java
  11. 2
      dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java
  12. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java
  13. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java
  14. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java

6
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<UUID, Integer> 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());

9
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<NotificationDeliveryMethod> deliveryMethods = ctx.getDeliveryMethods();
List<ListenableFuture<Void>> 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);
}

4
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)

4
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)

12
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<User, NotificationApiWsClient> sessions = new HashMap<>();
List<NotificationTargetId> 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);
}

2
common/dao-api/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestService.java

@ -44,6 +44,6 @@ public interface NotificationRequestService {
PageData<NotificationRequest> findScheduledNotificationRequests(PageLink pageLink);
void updateNotificationRequestStats(TenantId tenantId, NotificationRequestId notificationRequestId, NotificationRequestStats stats);
void updateNotificationRequest(TenantId tenantId, NotificationRequestId requestId, NotificationRequestStatus requestStatus, NotificationRequestStats stats);
}

5
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<NotificationRequestId> impleme
private TenantId tenantId;
@NotNull
private List<NotificationTargetId> targets;
private List<UUID> targets;
@NotNull
private NotificationTemplateId templateId;

1
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
}

2
dao/src/main/java/org/thingsboard/server/dao/model/sql/NotificationRequestEntity.java

@ -109,7 +109,7 @@ public class NotificationRequestEntity extends BaseSqlEntity<NotificationRequest
notificationRequest.setId(new NotificationRequestId(id));
notificationRequest.setCreatedTime(createdTime);
notificationRequest.setTenantId(getTenantId(tenantId));
notificationRequest.setTargets(listFromString(targets, uuid -> 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));

4
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);
}

2
dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java

@ -41,7 +41,7 @@ public interface NotificationRequestDao extends Dao<NotificationRequest> {
PageData<NotificationRequest> 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);

4
dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java

@ -72,8 +72,8 @@ public class JpaNotificationRequestDao extends JpaAbstractDao<NotificationReques
}
@Override
public void updateStatsById(TenantId tenantId, NotificationRequestId notificationRequestId, NotificationRequestStats stats) {
notificationRequestRepository.updateStatsById(notificationRequestId.getId(), JacksonUtil.valueToTree(stats));
public void updateById(TenantId tenantId, NotificationRequestId requestId, NotificationRequestStatus requestStatus, NotificationRequestStats stats) {
notificationRequestRepository.updateStatusAndStatsById(requestId.getId(), requestStatus, JacksonUtil.valueToTree(stats));
}
@Override

6
dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRequestRepository.java

@ -46,8 +46,10 @@ public interface NotificationRequestRepository extends JpaRepository<Notificatio
@Modifying
@Transactional
@Query("UPDATE NotificationRequestEntity r SET r.stats = :stats WHERE r.id = :id")
void updateStatsById(@Param("id") UUID id, @Param("stats") JsonNode stats);
@Query("UPDATE NotificationRequestEntity r SET r.status = :status, r.stats = :stats WHERE r.id = :id")
void updateStatusAndStatsById(@Param("id") UUID id,
@Param("status") NotificationRequestStatus status,
@Param("stats") JsonNode stats);
boolean existsByStatusAndTargetsContaining(NotificationRequestStatus status, String targetIdStr);

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java

@ -24,12 +24,13 @@ import org.thingsboard.server.common.data.notification.NotificationRequestConfig
import javax.validation.constraints.NotEmpty;
import javax.validation.constraints.NotNull;
import java.util.List;
import java.util.UUID;
@Data
public class TbNotificationNodeConfiguration implements NodeConfiguration<TbNotificationNodeConfiguration> {
@NotEmpty
private List<NotificationTargetId> targets;
private List<UUID> targets;
@NotNull
private NotificationTemplateId templateId;
private NotificationRequestConfig additionalConfig;

Loading…
Cancel
Save