diff --git a/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java b/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java index 5524255648..751b139bd3 100644 --- a/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java +++ b/application/src/main/java/org/thingsboard/server/controller/NotificationTemplateController.java @@ -17,6 +17,7 @@ package org.thingsboard.server.controller; import io.swagger.annotations.ApiOperation; import lombok.RequiredArgsConstructor; +import org.apache.commons.lang3.StringUtils; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.core.annotation.AuthenticationPrincipal; import org.springframework.web.bind.annotation.DeleteMapping; @@ -136,16 +137,20 @@ public class NotificationTemplateController extends BaseController { @GetMapping("/slack/conversations") @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") public List listSlackConversations(@RequestParam SlackConversationType type, + @RequestParam(required = false) String token, @AuthenticationPrincipal SecurityUser user) { // generic permission - NotificationSettings settings = notificationSettingsService.findNotificationSettings(user.getTenantId()); - SlackNotificationDeliveryMethodConfig slackConfig = (SlackNotificationDeliveryMethodConfig) - settings.getDeliveryMethodsConfigs().get(NotificationDeliveryMethod.SLACK); - if (slackConfig == null) { - throw new IllegalArgumentException("Slack is not configured"); + if (StringUtils.isEmpty(token)) { + NotificationSettings settings = notificationSettingsService.findNotificationSettings(user.getTenantId()); + SlackNotificationDeliveryMethodConfig slackConfig = (SlackNotificationDeliveryMethodConfig) + settings.getDeliveryMethodsConfigs().get(NotificationDeliveryMethod.SLACK); + if (slackConfig == null) { + throw new IllegalArgumentException("Slack is not configured"); + } + token = slackConfig.getBotToken(); } - return slackService.listConversations(user.getTenantId(), slackConfig.getBotToken(), type); + return slackService.listConversations(user.getTenantId(), token, type); } } 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 f8079d9a93..b4fba04212 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 @@ -364,6 +364,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple .deleted(true) .build()); } else if (notificationRequest.isScheduled()) { + // TODO: just forward to scheduler service clusterService.broadcastEntityStateChangeEvent(tenantId, notificationRequestId, ComponentLifecycleEvent.DELETED); } } diff --git a/application/src/main/java/org/thingsboard/server/service/slack/DefaultSlackService.java b/application/src/main/java/org/thingsboard/server/service/slack/DefaultSlackService.java index d6e351f51e..a36b12f0fc 100644 --- a/application/src/main/java/org/thingsboard/server/service/slack/DefaultSlackService.java +++ b/application/src/main/java/org/thingsboard/server/service/slack/DefaultSlackService.java @@ -155,7 +155,7 @@ public class DefaultSlackService implements SlackService { String neededScope = response.getNeeded(); error = "bot token scope '" + neededScope + "' is needed"; } - throw new RuntimeException("Failed to send message via Slack: " + error); + throw new RuntimeException("Slack API error: " + error); } return response; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java index b06a68e22c..60fce13560 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java @@ -316,20 +316,11 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene Set subscriptions = subscriptionsByEntityId.get(recipientId); if (subscriptions != null) { NotificationsSubscriptionUpdate subscriptionUpdate = new NotificationsSubscriptionUpdate(notificationUpdate); + log.trace("Handling notificationUpdate for user {}: {}", recipientId, notificationUpdate); subscriptions.stream() .filter(subscription -> subscription.getType() == TbSubscriptionType.NOTIFICATIONS || subscription.getType() == TbSubscriptionType.NOTIFICATIONS_COUNT) - .forEach(subscription -> { - if (serviceId.equals(subscription.getServiceId())) { - localSubscriptionService.onSubscriptionUpdate(subscription.getSessionId(), - subscription.getSubscriptionId(), subscriptionUpdate, TbCallback.EMPTY); - } else { - TopicPartitionInfo tpi = notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); - ToCoreNotificationMsg updateProto = TbSubscriptionUtils.notificationsSubUpdateToProto(subscription, subscriptionUpdate); - TbProtoQueueMsg queueMsg = new TbProtoQueueMsg<>(subscription.getEntityId().getId(), updateProto); - toCoreNotificationsProducer.send(tpi, queueMsg, null); - } - }); + .forEach(subscription -> onNotificationsSubUpdate(subscriptionUpdate, subscription)); } callback.onSuccess(); } @@ -341,21 +332,37 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene if (entityId.getEntityType() != EntityType.USER) { return; } + log.trace("Handling notificationRequestUpdate for user {}: {}", entityId, notificationRequestUpdate); subscriptions.forEach(subscription -> { if (subscription.getType() != TbSubscriptionType.NOTIFICATIONS && subscription.getType() != TbSubscriptionType.NOTIFICATIONS_COUNT) { return; } - if (!subscription.getTenantId().equals(tenantId) || !subscription.getServiceId().equals(serviceId)) { + if (!subscription.getTenantId().equals(tenantId)) { return; } - localSubscriptionService.onSubscriptionUpdate(subscription.getSessionId(), subscription.getSubscriptionId(), - subscriptionUpdate, TbCallback.EMPTY); + onNotificationsSubUpdate(subscriptionUpdate, subscription); }); }); callback.onSuccess(); } + private void onNotificationsSubUpdate(NotificationsSubscriptionUpdate subscriptionUpdate, TbSubscription subscription) { + if (serviceId.equals(subscription.getServiceId())) { + log.trace("[{}][{}][{}] Subscription session is managed by current service, forwarding to localSubscriptionService (update: {})", + subscription.getServiceId(), subscription.getEntityId(), subscription.getSessionId(), subscriptionUpdate); + localSubscriptionService.onSubscriptionUpdate(subscription.getSessionId(), + subscription.getSubscriptionId(), subscriptionUpdate, TbCallback.EMPTY); + } else { + log.trace("[{}][{}][{}] Subscription session is not managed by current service (update: {})", + subscription.getServiceId(), subscription.getEntityId(), subscription.getSessionId(), subscriptionUpdate); + TopicPartitionInfo tpi = notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); + ToCoreNotificationMsg updateProto = TbSubscriptionUtils.notificationsSubUpdateToProto(subscription, subscriptionUpdate); + TbProtoQueueMsg queueMsg = new TbProtoQueueMsg<>(subscription.getEntityId().getId(), updateProto); + toCoreNotificationsProducer.send(tpi, queueMsg, null); + } + } + @Override public void onAttributesDelete(TenantId tenantId, EntityId entityId, String scope, List keys, boolean notifyDevice, TbCallback callback) { onLocalTelemetrySubUpdate(entityId, 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 ec76319493..d85feb2f05 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 @@ -64,7 +64,7 @@ public class TbNotificationNode implements TbNode { NotificationRequest notificationRequest = NotificationRequest.builder() .tenantId(ctx.getTenantId()) .targets(config.getTargets()) - .templateId(new NotificationTemplateId(config.getTemplateId())) + .templateId(config.getTemplateId()) .info(notificationInfo) .additionalConfig(new NotificationRequestConfig()) .originatorEntityId(ctx.getSelf().getRuleChainId()) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java index f18f18313c..9508a6a848 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNodeConfiguration.java @@ -17,9 +17,7 @@ package org.thingsboard.rule.engine.notification; import lombok.Data; import org.thingsboard.rule.engine.api.NodeConfiguration; -import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.NotificationTemplateId; -import org.thingsboard.server.common.data.notification.NotificationRequestConfig; import javax.validation.constraints.NotEmpty; import javax.validation.constraints.NotNull; @@ -32,7 +30,7 @@ public class TbNotificationNodeConfiguration implements NodeConfiguration targets; @NotNull - private UUID templateId; + private NotificationTemplateId templateId; @Override public TbNotificationNodeConfiguration defaultConfiguration() { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java index 24b6000b72..544a864931 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java @@ -29,7 +29,7 @@ import java.util.concurrent.ExecutionException; @RuleNode( type = ComponentType.EXTERNAL, - name = "send to Slack", + name = "send to slack", configClazz = TbSlackNodeConfiguration.class, nodeDescription = "Send message via Slack", nodeDetails = "Sends message to a Slack channel or user", diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index e58ada965c..c64ca941f2 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1399,10 +1399,10 @@ "create-new-device-profile": "Create a new one!", "mqtt-device-topic-filters": "MQTT device topic filters", "mqtt-device-topic-filters-unique": "MQTT device topic filters need to be unique.", - "mqtt-device-topic-filters-spark-plug": "MQTT device topic filters SparkPlug.", - "mqtt-device-topic-filters-spark-plug-hint": "Default - telemetry. Example: namespace/group_id/message_type/edge_node_id/[device_id].", - "mqtt-device-topic-filters-spark-plug-attribute-metric-names": "SparkPlug attributes metric names", - "mqtt-device-topic-filters-spark-plug-attribute-metric-names-hint": "Names of SparkPlug metrics that will be stored as device attributes. All other metrics will be stored as device telemetry", + "mqtt-device-topic-filters-spark-plug": "MQTT Sparkplug B Edge of Network (EoN) node.", + "mqtt-device-topic-filters-spark-plug-hint": "Allow connections from EoN nodes with Sparkplug B payload and topic format.", + "mqtt-device-topic-filters-spark-plug-attribute-metric-names": "SparkPlug metrics to store as attributes.", + "mqtt-device-topic-filters-spark-plug-attribute-metric-names-hint": "Names of SparkPlug metrics that will be stored as device attributes. All other metrics will be stored as device telemetry.", "mqtt-device-payload-type": "MQTT device payload", "mqtt-device-payload-type-json": "JSON", "mqtt-device-payload-type-proto": "Protobuf",