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 b4ba871edf..ef4fd02e8b 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.notification; +import com.google.common.util.concurrent.FutureCallback; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -79,7 +80,6 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.UUID; -import java.util.function.Consumer; import java.util.stream.Collectors; @Service @@ -101,7 +101,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple private Map channels; @Override - public NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest request, Consumer callback) { + public NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest request, FutureCallback callback) { if (request.getRuleId() == null) { if (!rateLimitService.checkRateLimit(LimitedApi.NOTIFICATION_REQUESTS, tenantId)) { throw new TbRateLimitsException(EntityType.TENANT); @@ -200,7 +200,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple } } - private void processNotificationRequestAsync(NotificationProcessingContext ctx, List targets, Consumer callback) { + private void processNotificationRequestAsync(NotificationProcessingContext ctx, List targets, FutureCallback callback) { notificationExecutor.submit(() -> { NotificationRequestId requestId = ctx.getRequest().getId(); for (NotificationTarget target : targets) { @@ -208,33 +208,39 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple processForTarget(target, ctx); } catch (Exception e) { log.error("[{}] Failed to process notification request for target {}", requestId, target.getId(), e); + ctx.getStats().setError(e.getMessage()); + updateRequestStats(ctx, requestId, ctx.getStats()); + + if (callback != null) { + callback.onFailure(e); + } + return; } } log.debug("[{}] Notification request processing is finished", requestId); NotificationRequestStats stats = ctx.getStats(); - try { - notificationRequestService.updateNotificationRequest(ctx.getTenantId(), requestId, NotificationRequestStatus.SENT, stats); - } catch (Exception e) { - log.error("[{}] Failed to update stats for notification request", requestId, e); - } - + updateRequestStats(ctx, requestId, stats); if (callback != null) { - try { - callback.accept(stats); - } catch (Exception e) { - log.error("Failed to process callback for notification request {}", requestId, e); - } + callback.onSuccess(stats); } }); } + private void updateRequestStats(NotificationProcessingContext ctx, NotificationRequestId requestId, NotificationRequestStats stats) { + try { + notificationRequestService.updateNotificationRequest(ctx.getTenantId(), requestId, NotificationRequestStatus.SENT, stats); + } catch (Exception e) { + log.error("[{}] Failed to update stats for notification request", requestId, e); + } + } + private void processForTarget(NotificationTarget target, NotificationProcessingContext ctx) { Iterable recipients; switch (target.getConfiguration().getType()) { case PLATFORM_USERS: { PlatformUsersNotificationTargetConfig targetConfig = (PlatformUsersNotificationTargetConfig) target.getConfiguration(); - if (targetConfig.getUsersFilter().getType().isForRules()) { + if (targetConfig.getUsersFilter().getType().isForRules() && ctx.getRequest().getInfo() instanceof RuleOriginatedNotificationInfo) { recipients = new PageDataIterable<>(pageLink -> { return notificationTargetService.findRecipientsForRuleNotificationTargetConfig(ctx.getTenantId(), targetConfig, (RuleOriginatedNotificationInfo) ctx.getRequest().getInfo(), pageLink); }, 500); 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 ba0b4e4763..379f7648ad 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.notification; +import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.assertj.core.data.Offset; @@ -709,7 +710,17 @@ public class NotificationApiTest extends AbstractNotificationApiTest { private NotificationRequestStats submitNotificationRequestAndWait(NotificationRequest notificationRequest) throws Exception { SettableFuture future = SettableFuture.create(); - notificationCenter.processNotificationRequest(notificationRequest.getTenantId(), notificationRequest, future::set); + notificationCenter.processNotificationRequest(notificationRequest.getTenantId(), notificationRequest, new FutureCallback<>() { + @Override + public void onSuccess(NotificationRequestStats result) { + future.set(result); + } + + @Override + public void onFailure(Throwable t) { + future.setException(t); + } + }); return future.get(30, TimeUnit.SECONDS); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RuleEngineOriginatedNotificationInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RuleEngineOriginatedNotificationInfo.java index adf1e3a0f6..f37c145f2c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RuleEngineOriginatedNotificationInfo.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/RuleEngineOriginatedNotificationInfo.java @@ -19,6 +19,7 @@ import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import java.util.HashMap; @@ -28,9 +29,10 @@ import java.util.Map; @AllArgsConstructor @NoArgsConstructor @Builder -public class RuleEngineOriginatedNotificationInfo implements NotificationInfo { +public class RuleEngineOriginatedNotificationInfo implements RuleOriginatedNotificationInfo { private EntityId msgOriginator; + private CustomerId msgCustomerId; private String msgType; private Map msgMetadata; private Map msgData; @@ -43,6 +45,7 @@ public class RuleEngineOriginatedNotificationInfo implements NotificationInfo { templateData.put("originatorType", msgOriginator.getEntityType().getNormalName()); templateData.put("originatorId", msgOriginator.getId().toString()); templateData.put("msgType", msgType); + templateData.put("customerId", msgCustomerId != null ? msgCustomerId.getId().toString() : ""); return templateData; } @@ -51,4 +54,9 @@ public class RuleEngineOriginatedNotificationInfo implements NotificationInfo { return msgOriginator; } + @Override + public CustomerId getAffectedCustomerId() { + return msgCustomerId; + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java index 97a7ee0e27..f313fb6e9f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java @@ -150,8 +150,9 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl return userService.findAllUsers(pageLink); } } + default: + throw new IllegalArgumentException("Recipient type not supported"); } - return new PageData<>(); } @Override @@ -178,6 +179,8 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl return userService.findTenantAdmins(affectedTenantId, pageLink); } break; + default: + throw new IllegalArgumentException("Recipient type not supported"); } return new PageData<>(); } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java index 772a861dbd..7fd18c4afd 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.api; +import com.google.common.util.concurrent.FutureCallback; import org.thingsboard.server.common.data.id.NotificationId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; @@ -26,11 +27,10 @@ import org.thingsboard.server.common.data.notification.targets.platform.UsersFil import org.thingsboard.server.common.data.notification.template.NotificationTemplate; import java.util.Set; -import java.util.function.Consumer; public interface NotificationCenter { - NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest, Consumer callback); + NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest, FutureCallback callback); void sendGeneralWebNotification(TenantId tenantId, UsersFilter recipients, NotificationTemplate template); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java index 5f1bea2bdb..1db5a7fb11 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.notification; +import com.google.common.util.concurrent.FutureCallback; import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; @@ -23,8 +24,10 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.external.TbAbstractExternalNode; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.notification.NotificationRequestConfig; +import org.thingsboard.server.common.data.notification.NotificationRequestStats; import org.thingsboard.server.common.data.notification.info.RuleEngineOriginatedNotificationInfo; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -56,6 +59,8 @@ public class TbNotificationNode extends TbAbstractExternalNode { public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException { RuleEngineOriginatedNotificationInfo notificationInfo = RuleEngineOriginatedNotificationInfo.builder() .msgOriginator(msg.getOriginator()) + .msgCustomerId(msg.getOriginator().getEntityType() == EntityType.CUSTOMER + && msg.getOriginator().equals(msg.getCustomerId()) ? null : msg.getCustomerId()) .msgMetadata(msg.getMetaData().getData()) .msgData(JacksonUtil.toFlatMap(JacksonUtil.toJsonNode(msg.getData()))) .msgType(msg.getType()) @@ -72,15 +77,23 @@ public class TbNotificationNode extends TbAbstractExternalNode { var tbMsg = ackIfNeeded(ctx, msg); - DonAsynchron.withCallback(ctx.getNotificationExecutor().executeAsync(() -> - ctx.getNotificationCenter().processNotificationRequest(ctx.getTenantId(), notificationRequest, stats -> { - TbMsgMetaData metaData = tbMsg.getMetaData().copy(); - metaData.putValue("notificationRequestResult", JacksonUtil.toString(stats)); - tellSuccess(ctx, TbMsg.transformMsgMetadata(tbMsg, metaData)); - })), - r -> { - }, - e -> tellFailure(ctx, tbMsg, e)); + var callback = new FutureCallback() { + @Override + public void onSuccess(NotificationRequestStats stats) { + TbMsgMetaData metaData = tbMsg.getMetaData().copy(); + metaData.putValue("notificationRequestResult", JacksonUtil.toString(stats)); + tellSuccess(ctx, TbMsg.transformMsgMetadata(tbMsg, metaData)); + } + + @Override + public void onFailure(Throwable e) { + tellFailure(ctx, tbMsg, e); + } + }; + + var future = ctx.getNotificationExecutor().executeAsync(() -> + ctx.getNotificationCenter().processNotificationRequest(ctx.getTenantId(), notificationRequest, callback)); + DonAsynchron.withCallback(future, r -> {}, callback::onFailure); } } diff --git a/ui-ngx/src/app/shared/models/notification.models.ts b/ui-ngx/src/app/shared/models/notification.models.ts index c47be6f312..22dcb66042 100644 --- a/ui-ngx/src/app/shared/models/notification.models.ts +++ b/ui-ngx/src/app/shared/models/notification.models.ts @@ -442,14 +442,12 @@ export const NotificationTargetConfigTypeInfoMap = new Map