From 67656a275706857aa79147b048daf1fc95152629 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 2 Jun 2023 14:50:23 +0300 Subject: [PATCH 1/5] Notifications deduplication --- .../server/actors/ActorSystemContext.java | 2 +- .../DefaultNotificationRuleProcessor.java | 43 +++++++----- .../src/main/resources/thingsboard.yml | 5 ++ .../notification/NotificationRuleApiTest.java | 2 +- .../notification/rule/NotificationRule.java | 10 +++ .../NotificationRuleTriggerConfig.java | 6 ++ .../trigger/NotificationRuleTriggerType.java | 15 +++-- .../settings/TriggerTypeConfig.java | 14 ++-- .../NotificationRuleProcessor.java | 2 +- .../trigger/NewPlatformVersionTrigger.java | 17 +++++ .../trigger/NotificationRuleTrigger.java | 14 ++++ .../RemoteNotificationRuleProcessor.java | 65 ++++++++++++++++++- .../src/main/resources/tb-vc-executor.yml | 9 ++- .../src/main/resources/tb-coap-transport.yml | 7 ++ .../src/main/resources/tb-http-transport.yml | 7 ++ .../src/main/resources/tb-lwm2m-transport.yml | 7 ++ .../src/main/resources/tb-mqtt-transport.yml | 7 ++ .../src/main/resources/tb-snmp-transport.yml | 7 ++ 18 files changed, 201 insertions(+), 38 deletions(-) rename application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RuleEngineMsgNotificationRuleTriggerProcessor.java => common/data/src/main/java/org/thingsboard/server/common/data/notification/settings/TriggerTypeConfig.java (56%) rename common/{queue/src/main/java/org/thingsboard/server/queue => message/src/main/java/org/thingsboard/server/common/msg}/notification/NotificationRuleProcessor.java (93%) diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index e0190bb388..e04c7052db 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -90,7 +90,7 @@ import org.thingsboard.server.dao.widget.WidgetsBundleService; import org.thingsboard.server.queue.discovery.DiscoveryService; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.component.ComponentDiscoveryService; diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java index 9c57b3cc77..297ccc4007 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java @@ -66,6 +66,7 @@ import java.util.stream.Collectors; @Service @RequiredArgsConstructor +@ConfigurationProperties(prefix = "notification-system.rules") @Slf4j @SuppressWarnings({"rawtypes", "unchecked"}) public class DefaultNotificationRuleProcessor implements NotificationRuleProcessor { @@ -79,6 +80,8 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess private final NotificationExecutorService notificationExecutor; private final CacheManager cacheManager; private Cache sentNotifications; + @Setter + private Map triggerTypesConfigs; private final Map triggerProcessors = new EnumMap<>(NotificationRuleTriggerType.class); @@ -142,9 +145,6 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess log.debug("[{}] Rate limit for notification requests per rule was exceeded (rule '{}')", rule.getTenantId(), rule.getName()); return; } - if (trigger.getType().isDeduplicate() && alreadySent(rule.getId(), trigger)) { - return; - } NotificationInfo notificationInfo = constructNotificationInfo(trigger, triggerConfig); rule.getRecipientsConfig().getTargetsTable().forEach((delay, targets) -> { @@ -194,23 +194,34 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess return triggerProcessors.get(triggerConfig.getTriggerType()).constructNotificationInfo(trigger); } - private boolean alreadySent(NotificationRuleId ruleId, NotificationRuleTrigger trigger) { - String key = ruleId + "_" + trigger.getOriginatorEntityId(); - SentNotification sent = sentNotifications.get(key, SentNotification.class); - boolean alreadySent; - if (sent != null && sent.getTrigger().equals(trigger)) { - alreadySent = true; - log.debug("Notification for {} trigger was already sent, ignoring", trigger.getType()); - // updating cache anyway so that the value is not removed by ttl - } else { - alreadySent = false; - sent = new SentNotification(trigger); + private boolean alreadySent(NotificationRule rule, NotificationRuleTrigger trigger) { + String deduplicationKey = getDeduplicationKey(trigger, rule); + + boolean alreadySent = false; + Long lastSentTs = sentNotifications.get(deduplicationKey, Long.class); + if (lastSentTs != null) { + long deduplicationDuration = Optional.ofNullable(triggerTypesConfigs) + .map(triggerTypes -> triggerTypes.get(trigger.getType())) + .map(TriggerTypeConfig::getDeduplicationDuration) + .orElseGet(trigger::getDefaultDeduplicationDuration); + long passed = System.currentTimeMillis() - lastSentTs; + log.trace("Deduplicating trigger {} for rule '{}' by key '{}'. Deduplication duration: {} ms, passed: {} ms", + trigger.getType(), rule.getName(), deduplicationKey, deduplicationDuration, passed); + if (deduplicationDuration == 0 || passed <= deduplicationDuration) { + alreadySent = true; + } } - log.trace("[{}] Putting to sentNotifications cache: {}", ruleId, trigger); - sentNotifications.put(key, sent); + if (!alreadySent) { + lastSentTs = System.currentTimeMillis(); + } + sentNotifications.put(deduplicationKey, lastSentTs); return alreadySent; } + public static String getDeduplicationKey(NotificationRuleTrigger trigger, NotificationRule rule) { + return String.join("_", trigger.getDeduplicationKey(), rule.getDeduplicationKey()); + } + @EventListener(ComponentLifecycleMsg.class) public void onNotificationRuleDeleted(ComponentLifecycleMsg componentLifecycleMsg) { if (componentLifecycleMsg.getEvent() != ComponentLifecycleEvent.DELETED || diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 6b05df7feb..f780d42a0d 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1261,6 +1261,11 @@ vc: notification_system: thread_pool_size: "${TB_NOTIFICATION_SYSTEM_THREAD_POOL_SIZE:10}" + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" management: endpoints: diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java index 455b897fcd..044428e270 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java @@ -444,7 +444,7 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { } @Test - public void testNotificationsDeduplication() throws Exception { + public void testNotificationsDeduplication_newPlatformVersion() throws Exception { loginSysAdmin(); NewPlatformVersionNotificationRuleTriggerConfig triggerConfig = new NewPlatformVersionNotificationRuleTriggerConfig(); createNotificationRule(triggerConfig, "Test", "Test", createNotificationTarget(tenantAdminUserId).getId()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRule.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRule.java index 5ca8dff863..c88ade5bb1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRule.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRule.java @@ -36,6 +36,8 @@ import javax.validation.constraints.AssertTrue; import javax.validation.constraints.NotBlank; import javax.validation.constraints.NotNull; import java.io.Serializable; +import java.util.List; +import java.util.stream.Collectors; @Data @NoArgsConstructor @@ -84,4 +86,12 @@ public class NotificationRule extends BaseData implements Ha triggerType == recipientsConfig.getTriggerType(); } + @JsonIgnore + public String getDeduplicationKey() { + String targets = recipientsConfig.getTargetsTable().values().stream() + .flatMap(List::stream).sorted().map(Object::toString) + .collect(Collectors.joining(",")); + return String.join(":", targets, triggerConfig.getDeduplicationKey()); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java index 3406a802c7..b0eec28858 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerConfig.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.data.notification.rule.trigger; +import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonSubTypes.Type; @@ -39,4 +40,9 @@ public interface NotificationRuleTriggerConfig extends Serializable { NotificationRuleTriggerType getTriggerType(); + @JsonIgnore + default String getDeduplicationKey() { + return "#"; + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java index e79e7f1195..dff86f4ba1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/trigger/NotificationRuleTriggerType.java @@ -16,10 +16,8 @@ package org.thingsboard.server.common.data.notification.rule.trigger; import lombok.Getter; -import lombok.RequiredArgsConstructor; @Getter -@RequiredArgsConstructor public enum NotificationRuleTriggerType { ENTITY_ACTION, @@ -28,15 +26,18 @@ public enum NotificationRuleTriggerType { ALARM_ASSIGNMENT, DEVICE_ACTIVITY, RULE_ENGINE_COMPONENT_LIFECYCLE_EVENT, - NEW_PLATFORM_VERSION(false, true), - ENTITIES_LIMIT(false, false), - API_USAGE_LIMIT(false, false); + NEW_PLATFORM_VERSION(false), + ENTITIES_LIMIT(false), + API_USAGE_LIMIT(false); private final boolean tenantLevel; - private final boolean deduplicate; NotificationRuleTriggerType() { - this(true, false); + this(true); + } + + NotificationRuleTriggerType(boolean tenantLevel) { + this.tenantLevel = tenantLevel; } } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RuleEngineMsgNotificationRuleTriggerProcessor.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/settings/TriggerTypeConfig.java similarity index 56% rename from application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RuleEngineMsgNotificationRuleTriggerProcessor.java rename to common/data/src/main/java/org/thingsboard/server/common/data/notification/settings/TriggerTypeConfig.java index ac8f692f02..6991bb1277 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/RuleEngineMsgNotificationRuleTriggerProcessor.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/settings/TriggerTypeConfig.java @@ -13,15 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.notification.rule.trigger; +package org.thingsboard.server.common.data.notification.settings; -import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerConfig; -import org.thingsboard.server.common.msg.notification.trigger.RuleEngineMsgTrigger; - -import java.util.Set; - -public interface RuleEngineMsgNotificationRuleTriggerProcessor extends NotificationRuleTriggerProcessor { - - Set getSupportedMsgTypes(); +import lombok.Data; +@Data +public class TriggerTypeConfig { + private long deduplicationDuration; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/notification/NotificationRuleProcessor.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/NotificationRuleProcessor.java similarity index 93% rename from common/queue/src/main/java/org/thingsboard/server/queue/notification/NotificationRuleProcessor.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/notification/NotificationRuleProcessor.java index 1fa9d82863..380773c1b0 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/notification/NotificationRuleProcessor.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/NotificationRuleProcessor.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.queue.notification; +package org.thingsboard.server.common.msg.notification; import org.thingsboard.server.common.msg.notification.trigger.NotificationRuleTrigger; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NewPlatformVersionTrigger.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NewPlatformVersionTrigger.java index 0da2f0c57b..2bb88a708b 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NewPlatformVersionTrigger.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NewPlatformVersionTrigger.java @@ -43,4 +43,21 @@ public class NewPlatformVersionTrigger implements NotificationRuleTrigger { return TenantId.SYS_TENANT_ID; } + + @Override + public boolean deduplicate() { + return true; + } + + @Override + public String getDeduplicationKey() { + return String.join(":", NotificationRuleTrigger.super.getDeduplicationKey(), + updateInfo.getCurrentVersion(), updateInfo.getLatestVersion()); + } + + @Override + public long getDefaultDeduplicationDuration() { + return 0; + } + } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NotificationRuleTrigger.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NotificationRuleTrigger.java index b511062549..0cfc87a4df 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NotificationRuleTrigger.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/NotificationRuleTrigger.java @@ -29,4 +29,18 @@ public interface NotificationRuleTrigger extends Serializable { EntityId getOriginatorEntityId(); + + default boolean deduplicate() { + return false; + } + + default String getDeduplicationKey() { + EntityId originatorEntityId = getOriginatorEntityId(); + return String.join(":", getType().toString(), originatorEntityId.getEntityType().toString(), originatorEntityId.getId().toString()); + } + + default long getDefaultDeduplicationDuration() { + return 0; + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/notification/RemoteNotificationRuleProcessor.java b/common/queue/src/main/java/org/thingsboard/server/queue/notification/RemoteNotificationRuleProcessor.java index 3c3988389f..617fdbb50e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/notification/RemoteNotificationRuleProcessor.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/notification/RemoteNotificationRuleProcessor.java @@ -19,7 +19,12 @@ import com.google.protobuf.ByteString; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.stereotype.Service; +import org.springframework.util.ConcurrentReferenceHashMap; +import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; +import org.thingsboard.server.common.data.notification.settings.TriggerTypeConfig; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.common.msg.notification.trigger.NotificationRuleTrigger; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -30,10 +35,17 @@ import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.DataDecodingEncodingService; +import java.util.EnumMap; +import java.util.Map; import java.util.UUID; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.springframework.util.ConcurrentReferenceHashMap.ReferenceType.SOFT; @Service @ConditionalOnMissingBean(value = NotificationRuleProcessor.class, ignored = RemoteNotificationRuleProcessor.class) +@ConfigurationProperties(prefix = "notification-system.rules") @RequiredArgsConstructor @Slf4j public class RemoteNotificationRuleProcessor implements NotificationRuleProcessor { @@ -43,10 +55,16 @@ public class RemoteNotificationRuleProcessor implements NotificationRuleProcesso private final PartitionService partitionService; private final DataDecodingEncodingService encodingService; + private Map triggerTypesConfigs; + private final ConcurrentMap submittedTriggers = new ConcurrentReferenceHashMap<>(16, SOFT); + @Override public void process(NotificationRuleTrigger trigger) { + if (trigger.deduplicate() && alreadySubmitted(trigger)) { + return; + } try { - log.trace("Submitting notification rule trigger: {}", trigger); + log.debug("Submitting notification rule trigger: {}", trigger); TransportProtos.NotificationRuleProcessorMsg.Builder msg = TransportProtos.NotificationRuleProcessorMsg.newBuilder() .setTrigger(ByteString.copyFrom(encodingService.encode(trigger))); @@ -57,9 +75,52 @@ public class RemoteNotificationRuleProcessor implements NotificationRuleProcesso .setNotificationRuleProcessorMsg(msg) .build()), null); }); - } catch (Exception e) { + } catch (Throwable e) { log.error("Failed to submit notification rule trigger: {}", trigger, e); } } + private boolean alreadySubmitted(NotificationRuleTrigger trigger) { + String deduplicationKey = trigger.getDeduplicationKey(); + + AtomicBoolean alreadySubmitted = new AtomicBoolean(false); + submittedTriggers.compute(deduplicationKey, (key, lastSubmittedTs) -> { + long currentTs = System.currentTimeMillis(); + if (lastSubmittedTs == null) { + return currentTs; + } else { + long deduplicationDuration = getDeduplicationDuration(trigger); + long passed = currentTs - lastSubmittedTs; + if (deduplicationDuration == 0 || passed <= deduplicationDuration) { + log.trace("Notification rule trigger {} was already submitted {} ms ago, deduplication duration is {} ms. Key: '{}'", + trigger.getType(), passed, deduplicationDuration, deduplicationKey); + alreadySubmitted.set(true); + return lastSubmittedTs; + } else { + return currentTs; + } + } + }); + return alreadySubmitted.get(); + } + + private long getDeduplicationDuration(NotificationRuleTrigger trigger) { + if (triggerTypesConfigs == null) { + triggerTypesConfigs = new EnumMap<>(NotificationRuleTriggerType.class); + } + TriggerTypeConfig triggerTypeConfig = triggerTypesConfigs.computeIfAbsent(trigger.getType(), triggerType -> { + TriggerTypeConfig config = new TriggerTypeConfig(); + config.setDeduplicationDuration(trigger.getDefaultDeduplicationDuration()); + return config; + }); + return triggerTypeConfig.getDeduplicationDuration(); + } + + // set from ConfigurationProperties + public void setTriggerTypesConfigs(Map triggerTypesConfigs) { + if (triggerTypesConfigs != null) { + this.triggerTypesConfigs = new EnumMap<>(triggerTypesConfigs); + } + } + } diff --git a/msa/vc-executor/src/main/resources/tb-vc-executor.yml b/msa/vc-executor/src/main/resources/tb-vc-executor.yml index ca495c678f..75f9e09d3c 100644 --- a/msa/vc-executor/src/main/resources/tb-vc-executor.yml +++ b/msa/vc-executor/src/main/resources/tb-vc-executor.yml @@ -202,4 +202,11 @@ management: service: type: "${TB_SERVICE_TYPE:tb-vc-executor}" # Unique id for this service (autogenerated if empty) - id: "${TB_SERVICE_ID:}" \ No newline at end of file + id: "${TB_SERVICE_ID:}" + +notification_system: + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/coap/src/main/resources/tb-coap-transport.yml b/transport/coap/src/main/resources/tb-coap-transport.yml index 7ea553fe5c..8fe079859a 100644 --- a/transport/coap/src/main/resources/tb-coap-transport.yml +++ b/transport/coap/src/main/resources/tb-coap-transport.yml @@ -302,3 +302,10 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' + +notification_system: + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/http/src/main/resources/tb-http-transport.yml b/transport/http/src/main/resources/tb-http-transport.yml index 346ec48eae..f05db08643 100644 --- a/transport/http/src/main/resources/tb-http-transport.yml +++ b/transport/http/src/main/resources/tb-http-transport.yml @@ -287,3 +287,10 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' + +notification_system: + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml index 4e8167d89d..745a7d126a 100644 --- a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml +++ b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml @@ -369,3 +369,10 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' + +notification_system: + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml index 1e0b1ebcd4..3795c533a7 100644 --- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml +++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml @@ -317,3 +317,10 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' + +notification_system: + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 9f086bcbc5..11dcc96010 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -267,3 +267,10 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' + +notification_system: + rules: + trigger_types_configs: + RATE_LIMITS: + # In milliseconds, 4 hours by default + deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" From f5cd8a9a52989c188e6bad899d7ed9fbb285c41d Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 2 Jun 2023 14:57:00 +0300 Subject: [PATCH 2/5] Notification system - less async operations --- .../DefaultTbApiUsageStateService.java | 2 +- .../DefaultNotificationCenter.java | 103 +++++++----------- .../channels/EmailNotificationChannel.java | 20 ++-- .../channels/NotificationChannel.java | 3 +- .../channels/SlackNotificationChannel.java | 10 +- .../channels/SmsNotificationChannel.java | 13 +-- .../DefaultNotificationRuleProcessor.java | 62 +++++------ .../cache/DefaultNotificationRulesCache.java | 15 +-- .../queue/DefaultTbCoreConsumerService.java | 2 +- .../DefaultAlarmSubscriptionService.java | 2 +- .../service/update/DefaultUpdateService.java | 2 +- .../src/main/resources/thingsboard.yml | 6 +- .../NotificationRequestStats.java | 3 + .../common/data/util/CollectionsUtil.java | 16 ++- .../src/main/resources/tb-vc-executor.yml | 7 -- .../src/main/resources/tb-coap-transport.yml | 7 -- .../src/main/resources/tb-http-transport.yml | 7 -- .../src/main/resources/tb-lwm2m-transport.yml | 7 -- .../src/main/resources/tb-mqtt-transport.yml | 7 -- .../src/main/resources/tb-snmp-transport.yml | 7 -- 20 files changed, 100 insertions(+), 201 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 9e17641c69..7e2719598c 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -61,7 +61,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceM import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.service.apiusage.BaseApiUsageState.StatsCalculationResult; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.mail.MailExecutorService; 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 0cb9dc780b..7215776170 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,13 +15,10 @@ */ package org.thingsboard.server.service.notification; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; -import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.rule.engine.api.NotificationCenter; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.User; @@ -59,14 +56,12 @@ import org.thingsboard.server.dao.notification.NotificationService; import org.thingsboard.server.dao.notification.NotificationSettingsService; import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.notification.NotificationTemplateService; -import org.thingsboard.server.dao.user.UserService; +import org.thingsboard.server.dao.util.limits.LimitedApi; +import org.thingsboard.server.dao.util.limits.RateLimitService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.NotificationsTopicService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; -import org.thingsboard.server.dao.util.limits.LimitedApi; -import org.thingsboard.server.dao.util.limits.RateLimitService; -import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.NotificationExecutorService; import org.thingsboard.server.service.notification.channels.NotificationChannel; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; @@ -74,7 +69,6 @@ import org.thingsboard.server.service.telemetry.AbstractSubscriptionService; import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate; -import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; import java.util.List; @@ -87,7 +81,7 @@ import java.util.stream.Collectors; @Service @Slf4j @RequiredArgsConstructor -@SuppressWarnings({"UnstableApiUsage", "rawtypes"}) +@SuppressWarnings({"rawtypes"}) public class DefaultNotificationCenter extends AbstractSubscriptionService implements NotificationCenter, NotificationChannel { private final NotificationTargetService notificationTargetService; @@ -95,9 +89,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple private final NotificationService notificationService; private final NotificationTemplateService notificationTemplateService; private final NotificationSettingsService notificationSettingsService; - private final UserService userService; private final NotificationExecutorService notificationExecutor; - private final DbCallbackExecutorService dbCallbackExecutorService; private final NotificationsTopicService notificationsTopicService; private final TbQueueProducerProvider producerProvider; private final RateLimitService rateLimitService; @@ -172,37 +164,32 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple .build(); notificationExecutor.submit(() -> { - List> results = new ArrayList<>(); - for (NotificationTarget target : targets) { - List> result = processForTarget(target, ctx); - results.addAll(result); + processForTarget(target, ctx); + } + + NotificationRequestId requestId = ctx.getRequest().getId(); + log.debug("[{}] Notification request processing is finished", requestId); + NotificationRequestStats stats = ctx.getStats(); + try { + notificationRequestService.updateNotificationRequest(tenantId, requestId, NotificationRequestStatus.SENT, stats); + } catch (Exception e) { + log.error("[{}] Failed to update stats for notification request", requestId, e); } - Futures.whenAllComplete(results).run(() -> { - NotificationRequestId requestId = ctx.getRequest().getId(); - log.debug("[{}] Notification request processing is finished", requestId); - NotificationRequestStats stats = ctx.getStats(); + if (callback != null) { try { - notificationRequestService.updateNotificationRequest(tenantId, requestId, NotificationRequestStatus.SENT, stats); + callback.accept(stats); } catch (Exception e) { - log.error("[{}] Failed to update stats for notification request", requestId, e); - } - - if (callback != null) { - try { - callback.accept(stats); - } catch (Exception e) { - log.error("Failed to process callback for notification request {}", requestId, e); - } + log.error("Failed to process callback for notification request {}", requestId, e); } - }, dbCallbackExecutorService); + } }); return request; } - private List> processForTarget(NotificationTarget target, NotificationProcessingContext ctx) { + private void processForTarget(NotificationTarget target, NotificationProcessingContext ctx) { Iterable recipients; switch (target.getConfiguration().getType()) { case PLATFORM_USERS: { @@ -231,43 +218,35 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple Set deliveryMethods = new HashSet<>(ctx.getDeliveryMethods()); deliveryMethods.removeIf(deliveryMethod -> !target.getConfiguration().getType().getSupportedDeliveryMethods().contains(deliveryMethod)); log.debug("[{}] Processing notification request for {} target ({}) for delivery methods {}", ctx.getRequest().getId(), target.getConfiguration().getType(), target.getId(), deliveryMethods); + if (deliveryMethods.isEmpty()) { + return; + } - List> results = new ArrayList<>(); - if (!deliveryMethods.isEmpty()) { - for (NotificationRecipient recipient : recipients) { - for (NotificationDeliveryMethod deliveryMethod : deliveryMethods) { - ListenableFuture resultFuture = processForRecipient(deliveryMethod, recipient, ctx); - DonAsynchron.withCallback(resultFuture, result -> { - ctx.getStats().reportSent(deliveryMethod, recipient); - }, error -> { - ctx.getStats().reportError(deliveryMethod, error, recipient); - }); - results.add(resultFuture); + for (NotificationRecipient recipient : recipients) { + for (NotificationDeliveryMethod deliveryMethod : deliveryMethods) { + try { + processForRecipient(deliveryMethod, recipient, ctx); + ctx.getStats().reportSent(deliveryMethod, recipient); + } catch (Exception error) { + ctx.getStats().reportError(deliveryMethod, error, recipient); } } } - return results; } - private ListenableFuture processForRecipient(NotificationDeliveryMethod deliveryMethod, NotificationRecipient recipient, NotificationProcessingContext ctx) { + private void processForRecipient(NotificationDeliveryMethod deliveryMethod, NotificationRecipient recipient, NotificationProcessingContext ctx) throws Exception { if (ctx.getStats().contains(deliveryMethod, recipient.getId())) { - return Futures.immediateFailedFuture(new AlreadySentException()); + throw new AlreadySentException(); } - - DeliveryMethodNotificationTemplate processedTemplate; - try { - processedTemplate = ctx.getProcessedTemplate(deliveryMethod, recipient); - } catch (Exception e) { - return Futures.immediateFailedFuture(e); - } - NotificationChannel notificationChannel = channels.get(deliveryMethod); + DeliveryMethodNotificationTemplate processedTemplate = ctx.getProcessedTemplate(deliveryMethod, recipient); + log.trace("[{}] Sending {} notification for recipient {}", ctx.getRequest().getId(), deliveryMethod, recipient); - return notificationChannel.sendNotification(recipient, processedTemplate, ctx); + notificationChannel.sendNotification(recipient, processedTemplate, ctx); } @Override - public ListenableFuture sendNotification(User recipient, WebDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) { + public void sendNotification(User recipient, WebDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws Exception { NotificationRequest request = ctx.getRequest(); Notification notification = Notification.builder() .requestId(request.getId()) @@ -283,14 +262,14 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple notification = notificationService.saveNotification(recipient.getTenantId(), notification); } catch (Exception e) { log.error("Failed to create notification for recipient {}", recipient.getId(), e); - return Futures.immediateFailedFuture(e); + throw e; } NotificationUpdate update = NotificationUpdate.builder() .created(true) .notification(notification) .build(); - return onNotificationUpdate(recipient.getTenantId(), recipient.getId(), update); + onNotificationUpdate(recipient.getTenantId(), recipient.getId(), update); } @Override @@ -384,13 +363,11 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple clusterService.pushMsgToCore(tenantId, notificationRequestId, toCoreMsg, null); } - private ListenableFuture onNotificationUpdate(TenantId tenantId, UserId recipientId, NotificationUpdate update) { + private void onNotificationUpdate(TenantId tenantId, UserId recipientId, NotificationUpdate update) { log.trace("Submitting notification update for recipient {}: {}", recipientId, update); - return Futures.submit(() -> { - forwardToSubscriptionManagerService(tenantId, recipientId, subscriptionManagerService -> { - subscriptionManagerService.onNotificationUpdate(tenantId, recipientId, update, TbCallback.EMPTY); - }, () -> TbSubscriptionUtils.notificationUpdateToProto(tenantId, recipientId, update)); - }, wsCallBackExecutor); + forwardToSubscriptionManagerService(tenantId, recipientId, subscriptionManagerService -> { + subscriptionManagerService.onNotificationUpdate(tenantId, recipientId, update, TbCallback.EMPTY); + }, () -> TbSubscriptionUtils.notificationUpdateToProto(tenantId, recipientId, update)); } private void onNotificationRequestUpdate(TenantId tenantId, NotificationRequestUpdate update) { diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/EmailNotificationChannel.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/EmailNotificationChannel.java index 8b3c9551c2..8e9e8525e3 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/channels/EmailNotificationChannel.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/EmailNotificationChannel.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.notification.channels; -import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Component; import org.thingsboard.rule.engine.api.MailService; @@ -24,7 +23,6 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.template.EmailDeliveryMethodNotificationTemplate; -import org.thingsboard.server.service.mail.MailExecutorService; import org.thingsboard.server.service.notification.NotificationProcessingContext; @Component @@ -32,19 +30,15 @@ import org.thingsboard.server.service.notification.NotificationProcessingContext public class EmailNotificationChannel implements NotificationChannel { private final MailService mailService; - private final MailExecutorService executor; @Override - public ListenableFuture sendNotification(User recipient, EmailDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) { - return executor.submit(() -> { - mailService.send(recipient.getTenantId(), null, TbEmail.builder() - .to(recipient.getEmail()) - .subject(processedTemplate.getSubject()) - .body(processedTemplate.getBody()) - .html(true) - .build()); - return null; - }); + public void sendNotification(User recipient, EmailDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws Exception { + mailService.send(recipient.getTenantId(), null, TbEmail.builder() + .to(recipient.getEmail()) + .subject(processedTemplate.getSubject()) + .body(processedTemplate.getBody()) + .html(true) + .build()); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/NotificationChannel.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/NotificationChannel.java index 02fe6264d2..0275644765 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/channels/NotificationChannel.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/NotificationChannel.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.notification.channels; -import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.targets.NotificationRecipient; @@ -24,7 +23,7 @@ import org.thingsboard.server.service.notification.NotificationProcessingContext public interface NotificationChannel { - ListenableFuture sendNotification(R recipient, T processedTemplate, NotificationProcessingContext ctx); + void sendNotification(R recipient, T processedTemplate, NotificationProcessingContext ctx) throws Exception; void check(TenantId tenantId) throws Exception; diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java index 46afbd7270..25c6b9dcab 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/SlackNotificationChannel.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.notification.channels; -import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Component; import org.thingsboard.rule.engine.api.slack.SlackService; @@ -26,7 +25,6 @@ import org.thingsboard.server.common.data.notification.settings.SlackNotificatio import org.thingsboard.server.common.data.notification.targets.slack.SlackConversation; import org.thingsboard.server.common.data.notification.template.SlackDeliveryMethodNotificationTemplate; import org.thingsboard.server.dao.notification.NotificationSettingsService; -import org.thingsboard.server.service.executors.ExternalCallExecutorService; import org.thingsboard.server.service.notification.NotificationProcessingContext; @Component @@ -35,15 +33,11 @@ public class SlackNotificationChannel implements NotificationChannel sendNotification(SlackConversation conversation, SlackDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) { + public void sendNotification(SlackConversation conversation, SlackDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws Exception { SlackNotificationDeliveryMethodConfig config = ctx.getDeliveryMethodConfig(NotificationDeliveryMethod.SLACK); - return executor.submit(() -> { - slackService.sendMessage(ctx.getTenantId(), config.getBotToken(), conversation.getId(), processedTemplate.getBody()); - return null; - }); + slackService.sendMessage(ctx.getTenantId(), config.getBotToken(), conversation.getId(), processedTemplate.getBody()); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/notification/channels/SmsNotificationChannel.java b/application/src/main/java/org/thingsboard/server/service/notification/channels/SmsNotificationChannel.java index 3c44dabd4b..44b61d3bb3 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/channels/SmsNotificationChannel.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/channels/SmsNotificationChannel.java @@ -15,8 +15,6 @@ */ package org.thingsboard.server.service.notification.channels; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Component; @@ -26,26 +24,21 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.template.SmsDeliveryMethodNotificationTemplate; import org.thingsboard.server.service.notification.NotificationProcessingContext; -import org.thingsboard.server.service.sms.SmsExecutorService; @Component @RequiredArgsConstructor public class SmsNotificationChannel implements NotificationChannel { private final SmsService smsService; - private final SmsExecutorService executor; @Override - public ListenableFuture sendNotification(User recipient, SmsDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) { + public void sendNotification(User recipient, SmsDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws Exception { String phone = recipient.getPhone(); if (StringUtils.isBlank(phone)) { - return Futures.immediateFailedFuture(new RuntimeException("User does not have phone number")); + throw new RuntimeException("User does not have phone number"); } - return executor.submit(() -> { - smsService.sendSms(recipient.getTenantId(), recipient.getCustomerId(), new String[]{phone}, processedTemplate.getBody()); - return null; - }); + smsService.sendSms(recipient.getTenantId(), recipient.getCustomerId(), new String[]{phone}, processedTemplate.getBody()); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java index 297ccc4007..71a59026a2 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java @@ -15,10 +15,11 @@ */ package org.thingsboard.server.service.notification.rule; -import lombok.Data; import lombok.RequiredArgsConstructor; +import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.cache.Cache; import org.springframework.cache.CacheManager; import org.springframework.context.annotation.Lazy; @@ -38,29 +39,27 @@ import org.thingsboard.server.common.data.notification.info.NotificationInfo; import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; +import org.thingsboard.server.common.data.notification.settings.TriggerTypeConfig; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.common.msg.notification.trigger.NotificationRuleTrigger; -import org.thingsboard.server.common.msg.notification.trigger.RuleEngineMsgTrigger; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.dao.notification.NotificationRequestService; import org.thingsboard.server.dao.util.limits.LimitedApi; import org.thingsboard.server.dao.util.limits.RateLimitService; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; import org.thingsboard.server.service.executors.NotificationExecutorService; import org.thingsboard.server.service.notification.rule.cache.NotificationRulesCache; import org.thingsboard.server.service.notification.rule.trigger.NotificationRuleTriggerProcessor; -import org.thingsboard.server.service.notification.rule.trigger.RuleEngineMsgNotificationRuleTriggerProcessor; import javax.annotation.PostConstruct; -import java.io.Serializable; +import java.util.ArrayList; import java.util.Collection; import java.util.EnumMap; -import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Set; +import java.util.Optional; import java.util.UUID; import java.util.stream.Collectors; @@ -96,20 +95,27 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess @Override public void process(NotificationRuleTrigger trigger) { NotificationRuleTriggerType triggerType = trigger.getType(); - if (triggerType == null) return; TenantId tenantId = triggerType.isTenantLevel() ? trigger.getTenantId() : TenantId.SYS_TENANT_ID; try { - List rules = notificationRulesCache.getEnabled(tenantId, triggerType); - for (NotificationRule rule : rules) { - notificationExecutor.submit(() -> { + List enabledRules = notificationRulesCache.getEnabled(tenantId, triggerType); + if (enabledRules.isEmpty()) { + return; + } + if (trigger.deduplicate()) { + enabledRules = new ArrayList<>(enabledRules); + enabledRules.removeIf(rule -> alreadySent(rule, trigger)); + } + final List rules = enabledRules; + notificationExecutor.submit(() -> { + for (NotificationRule rule : rules) { try { processNotificationRule(rule, trigger); } catch (Throwable e) { log.error("Failed to process notification rule {} for trigger type {} with trigger object {}", rule.getId(), rule.getTriggerType(), trigger, e); } - }); - } + } + }); } catch (Throwable e) { log.error("Failed to process notification rules for trigger: {}", trigger, e); } @@ -172,14 +178,13 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess .ruleId(rule.getId()) .originatorEntityId(originatorEntityId) .build(); - notificationExecutor.submit(() -> { - try { - log.debug("Submitting notification request for rule '{}' with delay of {} sec to targets {}", rule.getName(), delayInSec, targets); - notificationCenter.processNotificationRequest(rule.getTenantId(), notificationRequest, null); - } catch (Exception e) { - log.error("Failed to process notification request for tenant {} for rule {}", rule.getTenantId(), rule.getId(), e); - } - }); + + try { + log.debug("Submitting notification request for rule '{}' with delay of {} sec to targets {}", rule.getName(), delayInSec, targets); + notificationCenter.processNotificationRequest(rule.getTenantId(), notificationRequest, null); + } catch (Exception e) { + log.error("Failed to process notification request for tenant {} for rule {}", rule.getTenantId(), rule.getId(), e); + } } private boolean matchesFilter(NotificationRuleTrigger trigger, NotificationRuleTriggerConfig triggerConfig) { @@ -243,24 +248,9 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess @Autowired public void setTriggerProcessors(Collection processors) { - Map ruleEngineMsgTypeToTriggerType = new HashMap<>(); processors.forEach(processor -> { triggerProcessors.put(processor.getTriggerType(), processor); - if (processor instanceof RuleEngineMsgNotificationRuleTriggerProcessor) { - Set supportedMsgTypes = ((RuleEngineMsgNotificationRuleTriggerProcessor) processor).getSupportedMsgTypes(); - supportedMsgTypes.forEach(supportedMsgType -> { - ruleEngineMsgTypeToTriggerType.put(supportedMsgType, processor.getTriggerType()); - }); - } }); - RuleEngineMsgTrigger.msgTypeToTriggerType = ruleEngineMsgTypeToTriggerType; - } - - @Data - private static class SentNotification implements Serializable { - private static final long serialVersionUID = 38973480405095422L; - - private final NotificationRuleTrigger trigger; } } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/cache/DefaultNotificationRulesCache.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/cache/DefaultNotificationRulesCache.java index a02b7bf1e4..a43ed59efa 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/cache/DefaultNotificationRulesCache.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/cache/DefaultNotificationRulesCache.java @@ -17,7 +17,6 @@ package org.thingsboard.server.service.notification.rule.cache; import com.github.benmanes.caffeine.cache.Cache; import com.github.benmanes.caffeine.cache.Caffeine; -import lombok.Data; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; @@ -50,7 +49,7 @@ public class DefaultNotificationRulesCache implements NotificationRulesCache { private int cacheMaxSize; @Value("${cache.notificationRules.timeToLiveInMinutes:30}") private int cacheValueTtl; - private Cache> cache; + private Cache> cache; private final ReadWriteLock lock = new ReentrantReadWriteLock(); @@ -96,21 +95,15 @@ public class DefaultNotificationRulesCache implements NotificationRulesCache { } } - private void evict(TenantId tenantId) { + public void evict(TenantId tenantId) { cache.invalidateAll(Arrays.stream(NotificationRuleTriggerType.values()) .map(triggerType -> key(tenantId, triggerType)) .collect(Collectors.toList())); log.trace("Evicted all notification rules for tenant {} from cache", tenantId); } - private static CacheKey key(TenantId tenantId, NotificationRuleTriggerType triggerType) { - return new CacheKey(tenantId, triggerType); - } - - @Data - private static class CacheKey { - private final TenantId tenantId; - private final NotificationRuleTriggerType triggerType; + private static String key(TenantId tenantId, NotificationRuleTriggerType triggerType) { + return tenantId + "_" + triggerType; } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 2bd0043046..182fce1c48 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -63,7 +63,7 @@ import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.DataDecodingEncodingService; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java index fec9d10a18..e7f2a5a66f 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java @@ -52,7 +52,7 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.stats.TbApiUsageReportClient; import org.thingsboard.server.dao.alarm.AlarmOperationResult; import org.thingsboard.server.dao.alarm.AlarmService; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.entitiy.alarm.TbAlarmCommentService; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; diff --git a/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java b/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java index ef066a73f6..587d86af32 100644 --- a/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java @@ -29,7 +29,7 @@ import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.UpdateMessage; import org.thingsboard.server.common.msg.notification.trigger.NewPlatformVersionTrigger; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbCoreComponent; diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index f780d42a0d..392587c869 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1263,9 +1263,9 @@ notification_system: thread_pool_size: "${TB_NOTIFICATION_SYSTEM_THREAD_POOL_SIZE:10}" rules: trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" + NEW_PLATFORM_VERSION: + # In milliseconds, infinitely by default + deduplication_duration: "${NEW_PLATFORM_VERSION_NOTIFICATION_RULE_DEDUPLICATION_DURATION:0}" management: endpoints: diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStats.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStats.java index 31e05a2cb3..619e1ad38f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStats.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/NotificationRequestStats.java @@ -62,6 +62,9 @@ public class NotificationRequestStats { return; } String errorMessage = error.getMessage(); + if (errorMessage == null) { + errorMessage = error.getClass().getSimpleName(); + } errors.computeIfAbsent(deliveryMethod, k -> new ConcurrentHashMap<>()).put(recipient.getTitle(), errorMessage); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java index d8d17613de..a0736ca50e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java @@ -16,7 +16,6 @@ package org.thingsboard.server.common.data.util; import java.util.Collection; -import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.Set; @@ -51,20 +50,19 @@ public class CollectionsUtil { } @SuppressWarnings("unchecked") - public static Map mapOf(Object... kvs) { - Map map = new HashMap<>(); + public static Map mapOf(T... kvs) { + if (kvs.length % 2 != 0) { + throw new IllegalArgumentException("Invalid number of parameters"); + } + Map map = new HashMap<>(); for (int i = 0; i < kvs.length; i += 2) { - K key = (K) kvs[i]; - V value = (V) kvs[i + 1]; + T key = kvs[i]; + T value = kvs[i + 1]; map.put(key, value); } return map; } - public static Map unmodifiableMapOf(Object... kvs) { - return Collections.unmodifiableMap(mapOf(kvs)); - } - public static boolean emptyOrContains(Collection collection, V element) { return isEmpty(collection) || collection.contains(element); } diff --git a/msa/vc-executor/src/main/resources/tb-vc-executor.yml b/msa/vc-executor/src/main/resources/tb-vc-executor.yml index 75f9e09d3c..9d5ad35388 100644 --- a/msa/vc-executor/src/main/resources/tb-vc-executor.yml +++ b/msa/vc-executor/src/main/resources/tb-vc-executor.yml @@ -203,10 +203,3 @@ service: type: "${TB_SERVICE_TYPE:tb-vc-executor}" # Unique id for this service (autogenerated if empty) id: "${TB_SERVICE_ID:}" - -notification_system: - rules: - trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/coap/src/main/resources/tb-coap-transport.yml b/transport/coap/src/main/resources/tb-coap-transport.yml index 8fe079859a..7ea553fe5c 100644 --- a/transport/coap/src/main/resources/tb-coap-transport.yml +++ b/transport/coap/src/main/resources/tb-coap-transport.yml @@ -302,10 +302,3 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' - -notification_system: - rules: - trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/http/src/main/resources/tb-http-transport.yml b/transport/http/src/main/resources/tb-http-transport.yml index f05db08643..346ec48eae 100644 --- a/transport/http/src/main/resources/tb-http-transport.yml +++ b/transport/http/src/main/resources/tb-http-transport.yml @@ -287,10 +287,3 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' - -notification_system: - rules: - trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml index 745a7d126a..4e8167d89d 100644 --- a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml +++ b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml @@ -369,10 +369,3 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' - -notification_system: - rules: - trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml index 3795c533a7..1e0b1ebcd4 100644 --- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml +++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml @@ -317,10 +317,3 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' - -notification_system: - rules: - trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 11dcc96010..9f086bcbc5 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -267,10 +267,3 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' - -notification_system: - rules: - trigger_types_configs: - RATE_LIMITS: - # In milliseconds, 4 hours by default - deduplication_duration: "${RATE_LIMITS_NOTIFICATION_RULE_DEDUPLICATION_DURATION:14400000}" From 64571eeff22f7f1e8352fb70231d6f58863337d5 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 2 Jun 2023 14:58:12 +0300 Subject: [PATCH 3/5] Don't use Rule Engine message as notification rule trigger --- .../server/actors/tenant/TenantActor.java | 5 -- .../controller/AlarmCommentController.java | 2 +- .../service/action/EntityActionService.java | 67 ++++++++++++++---- .../AlarmAssignmentTriggerProcessor.java | 39 ++++------- .../trigger/AlarmCommentTriggerProcessor.java | 70 +++++++++---------- .../DeviceActivityTriggerProcessor.java | 38 ++++------ .../state/DefaultDeviceStateService.java | 25 +++++-- .../state/DefaultDeviceStateServiceTest.java | 6 +- ...igger.java => AlarmAssignmentTrigger.java} | 22 +++--- .../trigger/AlarmCommentTrigger.java | 48 +++++++++++++ .../trigger/ApiUsageLimitTrigger.java | 5 -- .../trigger/DeviceActivityTrigger.java | 49 +++++++++++++ 12 files changed, 242 insertions(+), 134 deletions(-) rename common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/{RuleEngineMsgTrigger.java => AlarmAssignmentTrigger.java} (71%) create mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/AlarmCommentTrigger.java create mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/DeviceActivityTrigger.java diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index 84b7757846..11a5895fef 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -49,7 +49,6 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; -import org.thingsboard.server.common.msg.notification.trigger.RuleEngineMsgTrigger; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; @@ -211,10 +210,6 @@ public class TenantActor extends RuleChainManagerActor { log.trace("[{}] Ack message because Rule Engine is disabled", tenantId); tbMsg.getCallback().onSuccess(); } - systemContext.getNotificationRuleProcessor().process(RuleEngineMsgTrigger.builder() - .tenantId(tenantId) - .msg(tbMsg) - .build()); } private void onRuleChainMsg(RuleChainAwareMsg msg) { diff --git a/application/src/main/java/org/thingsboard/server/controller/AlarmCommentController.java b/application/src/main/java/org/thingsboard/server/controller/AlarmCommentController.java index 92b2cb923f..12195e83bc 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AlarmCommentController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AlarmCommentController.java @@ -77,7 +77,7 @@ public class AlarmCommentController extends BaseController { @PathVariable(ALARM_ID) String strAlarmId, @ApiParam(value = "A JSON value representing the comment.") @RequestBody AlarmComment alarmComment) throws ThingsboardException { checkParameter(ALARM_ID, strAlarmId); AlarmId alarmId = new AlarmId(toUUID(strAlarmId)); - Alarm alarm = checkAlarmId(alarmId, Operation.WRITE); + Alarm alarm = checkAlarmInfoId(alarmId, Operation.WRITE); alarmComment.setAlarmId(alarmId); return tbAlarmCommentService.saveAlarmComment(alarm, alarmComment, getCurrentUser()); } diff --git a/application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java b/application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java index 623b95b40e..99b508a12d 100644 --- a/application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java +++ b/application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java @@ -28,7 +28,9 @@ import org.thingsboard.server.common.data.HasName; import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmComment; +import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.CustomerId; @@ -40,10 +42,12 @@ import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.trigger.AlarmAssignmentTrigger; +import org.thingsboard.server.common.msg.notification.trigger.AlarmCommentTrigger; import org.thingsboard.server.common.msg.notification.trigger.EntitiesLimitTrigger; import org.thingsboard.server.common.msg.notification.trigger.EntityActionTrigger; import org.thingsboard.server.dao.audit.AuditLogService; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; import java.util.List; import java.util.Map; @@ -241,19 +245,7 @@ public class EntityActionService { } } if (tenantId != null && !tenantId.isSysTenantId()) { - if (actionType == ActionType.ADDED) { - notificationRuleProcessor.process(EntitiesLimitTrigger.builder() - .tenantId(tenantId) - .entityType(entityId.getEntityType()) - .build()); - } - notificationRuleProcessor.process(EntityActionTrigger.builder() - .tenantId(tenantId) - .entityId(entityId) - .entity(entity) - .actionType(actionType) - .user(user) - .build()); + processNotificationRules(tenantId, entityId, entity, actionType, user, additionalInfo); } TbMsg tbMsg = TbMsg.newMsg(msgType, entityId, customerId, metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode)); tbClusterService.pushMsgToRuleEngine(tenantId, entityId, tbMsg, null); @@ -263,6 +255,53 @@ public class EntityActionService { } } + private void processNotificationRules(TenantId tenantId, EntityId entityId, HasName entity, ActionType actionType, User user, Object... additionalInfo) { + switch (actionType) { + case ADDED: + notificationRuleProcessor.process(EntitiesLimitTrigger.builder() + .tenantId(tenantId) + .entityType(entityId.getEntityType()) + .build()); + case UPDATED: + case DELETED: + notificationRuleProcessor.process(EntityActionTrigger.builder() + .tenantId(tenantId) + .entityId(entityId) + .entity(entity) + .actionType(actionType) + .user(user) + .build()); + break; + case ALARM_ASSIGNED: + case ALARM_UNASSIGNED: + if (!(entity instanceof AlarmInfo)) { // should not normally happen + log.warn("Invalid alarm assignment event: entity is not instance of AlarmInfo"); + break; + } + notificationRuleProcessor.process(AlarmAssignmentTrigger.builder() + .tenantId(tenantId) + .alarmInfo((AlarmInfo) entity) + .actionType(actionType) + .user(user) + .build()); + break; + case ADDED_COMMENT: + case UPDATED_COMMENT: + if (!(entity instanceof Alarm)) { // should not normally happen + log.warn("Invalid alarm comment event: entity is not instance of Alarm"); + break; + } + notificationRuleProcessor.process(AlarmCommentTrigger.builder() + .tenantId(tenantId) + .comment(extractParameter(AlarmComment.class, 0, additionalInfo)) + .alarm((Alarm) entity) + .actionType(actionType) + .user(user) + .build()); + break; + } + } + public void logEntityAction(User user, I entityId, E entity, CustomerId customerId, ActionType actionType, Exception e, Object... additionalInfo) { if (customerId == null || customerId.isNullUid()) { diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmAssignmentTriggerProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmAssignmentTriggerProcessor.java index ed72f6be58..eca258aecc 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmAssignmentTriggerProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmAssignmentTriggerProcessor.java @@ -16,52 +16,48 @@ package org.thingsboard.server.service.notification.rule.trigger; import org.springframework.stereotype.Service; -import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmAssignee; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmStatusFilter; +import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.notification.info.AlarmAssignmentNotificationInfo; import org.thingsboard.server.common.data.notification.info.RuleOriginatedNotificationInfo; import org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentNotificationRuleTriggerConfig.Action; import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; -import org.thingsboard.server.common.msg.notification.trigger.RuleEngineMsgTrigger; - -import java.util.Set; +import org.thingsboard.server.common.msg.notification.trigger.AlarmAssignmentTrigger; import static org.apache.commons.collections.CollectionUtils.isEmpty; import static org.thingsboard.server.common.data.util.CollectionsUtil.emptyOrContains; @Service -public class AlarmAssignmentTriggerProcessor implements RuleEngineMsgNotificationRuleTriggerProcessor { +public class AlarmAssignmentTriggerProcessor implements NotificationRuleTriggerProcessor { @Override - public boolean matchesFilter(RuleEngineMsgTrigger trigger, AlarmAssignmentNotificationRuleTriggerConfig triggerConfig) { - Action action = trigger.getMsg().getType().equals(DataConstants.ALARM_ASSIGNED) ? Action.ASSIGNED : Action.UNASSIGNED; + public boolean matchesFilter(AlarmAssignmentTrigger trigger, AlarmAssignmentNotificationRuleTriggerConfig triggerConfig) { + Action action = trigger.getActionType() == ActionType.ALARM_ASSIGNED ? Action.ASSIGNED : Action.UNASSIGNED; if (!triggerConfig.getNotifyOn().contains(action)) { return false; } - Alarm alarm = JacksonUtil.fromString(trigger.getMsg().getData(), Alarm.class); - return emptyOrContains(triggerConfig.getAlarmTypes(), alarm.getType()) && - emptyOrContains(triggerConfig.getAlarmSeverities(), alarm.getSeverity()) && - (isEmpty(triggerConfig.getAlarmStatuses()) || AlarmStatusFilter.from(triggerConfig.getAlarmStatuses()).matches(alarm)); + AlarmInfo alarmInfo = trigger.getAlarmInfo(); + return emptyOrContains(triggerConfig.getAlarmTypes(), alarmInfo.getType()) && + emptyOrContains(triggerConfig.getAlarmSeverities(), alarmInfo.getSeverity()) && + (isEmpty(triggerConfig.getAlarmStatuses()) || AlarmStatusFilter.from(triggerConfig.getAlarmStatuses()).matches(alarmInfo)); } @Override - public RuleOriginatedNotificationInfo constructNotificationInfo(RuleEngineMsgTrigger trigger) { - AlarmInfo alarmInfo = JacksonUtil.fromString(trigger.getMsg().getData(), AlarmInfo.class); + public RuleOriginatedNotificationInfo constructNotificationInfo(AlarmAssignmentTrigger trigger) { + AlarmInfo alarmInfo = trigger.getAlarmInfo(); AlarmAssignee assignee = alarmInfo.getAssignee(); return AlarmAssignmentNotificationInfo.builder() - .action(trigger.getMsg().getType().equals(DataConstants.ALARM_ASSIGNED) ? "assigned" : "unassigned") + .action(trigger.getActionType() == ActionType.ALARM_ASSIGNED ? "assigned" : "unassigned") .assigneeFirstName(assignee != null ? assignee.getFirstName() : null) .assigneeLastName(assignee != null ? assignee.getLastName() : null) .assigneeEmail(assignee != null ? assignee.getEmail() : null) .assigneeId(assignee != null ? assignee.getId() : null) - .userEmail(trigger.getMsg().getMetaData().getValue("userEmail")) - .userFirstName(trigger.getMsg().getMetaData().getValue("userFirstName")) - .userLastName(trigger.getMsg().getMetaData().getValue("userLastName")) + .userEmail(trigger.getUser().getEmail()) + .userFirstName(trigger.getUser().getFirstName()) + .userLastName(trigger.getUser().getLastName()) .alarmId(alarmInfo.getUuidId()) .alarmType(alarmInfo.getType()) .alarmOriginator(alarmInfo.getOriginator()) @@ -77,9 +73,4 @@ public class AlarmAssignmentTriggerProcessor implements RuleEngineMsgNotificatio return NotificationRuleTriggerType.ALARM_ASSIGNMENT; } - @Override - public Set getSupportedMsgTypes() { - return Set.of(DataConstants.ALARM_ASSIGNED, DataConstants.ALARM_UNASSIGNED); - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmCommentTriggerProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmCommentTriggerProcessor.java index cf71718f83..70024d6e09 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmCommentTriggerProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/AlarmCommentTriggerProcessor.java @@ -15,68 +15,67 @@ */ package org.thingsboard.server.service.notification.rule.trigger; +import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; -import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.alarm.Alarm; -import org.thingsboard.server.common.data.alarm.AlarmComment; import org.thingsboard.server.common.data.alarm.AlarmCommentType; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmStatusFilter; +import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.notification.info.AlarmCommentNotificationInfo; import org.thingsboard.server.common.data.notification.info.RuleOriginatedNotificationInfo; import org.thingsboard.server.common.data.notification.rule.trigger.AlarmCommentNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.notification.trigger.RuleEngineMsgTrigger; - -import java.util.Set; +import org.thingsboard.server.common.msg.notification.trigger.AlarmCommentTrigger; +import org.thingsboard.server.dao.entity.EntityService; import static org.apache.commons.collections.CollectionUtils.isEmpty; import static org.thingsboard.server.common.data.util.CollectionsUtil.emptyOrContains; @Service -public class AlarmCommentTriggerProcessor implements RuleEngineMsgNotificationRuleTriggerProcessor { +@RequiredArgsConstructor +public class AlarmCommentTriggerProcessor implements NotificationRuleTriggerProcessor { + + private final EntityService entityService; @Override - public boolean matchesFilter(RuleEngineMsgTrigger trigger, AlarmCommentNotificationRuleTriggerConfig triggerConfig) { - TbMsg msg = trigger.getMsg(); - if (msg.getMetaData().getValue("comment") == null) { - return false; - } - if (msg.getType().equals(DataConstants.COMMENT_UPDATED) && !triggerConfig.isNotifyOnCommentUpdate()) { + public boolean matchesFilter(AlarmCommentTrigger trigger, AlarmCommentNotificationRuleTriggerConfig triggerConfig) { + if (trigger.getActionType() == ActionType.UPDATED_COMMENT && !triggerConfig.isNotifyOnCommentUpdate()) { return false; } if (triggerConfig.isOnlyUserComments()) { - AlarmComment comment = JacksonUtil.fromString(msg.getMetaData().getValue("comment"), AlarmComment.class); - if (comment.getType() == AlarmCommentType.SYSTEM) { + if (trigger.getComment().getType() == AlarmCommentType.SYSTEM) { return false; } } - Alarm alarm = JacksonUtil.fromString(msg.getData(), Alarm.class); + Alarm alarm = trigger.getAlarm(); return emptyOrContains(triggerConfig.getAlarmTypes(), alarm.getType()) && emptyOrContains(triggerConfig.getAlarmSeverities(), alarm.getSeverity()) && (isEmpty(triggerConfig.getAlarmStatuses()) || AlarmStatusFilter.from(triggerConfig.getAlarmStatuses()).matches(alarm)); } @Override - public RuleOriginatedNotificationInfo constructNotificationInfo(RuleEngineMsgTrigger trigger) { - TbMsg msg = trigger.getMsg(); - AlarmComment comment = JacksonUtil.fromString(msg.getMetaData().getValue("comment"), AlarmComment.class); - AlarmInfo alarmInfo = JacksonUtil.fromString(msg.getData(), AlarmInfo.class); + public RuleOriginatedNotificationInfo constructNotificationInfo(AlarmCommentTrigger trigger) { + Alarm alarm = trigger.getAlarm(); + String originatorName; + if (alarm instanceof AlarmInfo) { + originatorName = ((AlarmInfo) alarm).getOriginatorName(); + } else { + originatorName = entityService.fetchEntityName(trigger.getTenantId(), alarm.getOriginator()).orElse(""); + } return AlarmCommentNotificationInfo.builder() - .comment(comment.getComment().get("text").asText()) - .action(msg.getType().equals(DataConstants.COMMENT_CREATED) ? "added" : "updated") - .userEmail(msg.getMetaData().getValue("userEmail")) - .userFirstName(msg.getMetaData().getValue("userFirstName")) - .userLastName(msg.getMetaData().getValue("userLastName")) - .alarmId(alarmInfo.getUuidId()) - .alarmType(alarmInfo.getType()) - .alarmOriginator(alarmInfo.getOriginator()) - .alarmOriginatorName(alarmInfo.getOriginatorName()) - .alarmSeverity(alarmInfo.getSeverity()) - .alarmStatus(alarmInfo.getStatus()) - .alarmCustomerId(alarmInfo.getCustomerId()) + .comment(trigger.getComment().getComment().get("text").asText()) + .action(trigger.getActionType() == ActionType.ADDED_COMMENT ? "added" : "updated") + .userEmail(trigger.getUser().getEmail()) + .userFirstName(trigger.getUser().getFirstName()) + .userLastName(trigger.getUser().getLastName()) + .alarmId(alarm.getUuidId()) + .alarmType(alarm.getType()) + .alarmOriginator(alarm.getOriginator()) + .alarmOriginatorName(originatorName) + .alarmSeverity(alarm.getSeverity()) + .alarmStatus(alarm.getStatus()) + .alarmCustomerId(alarm.getCustomerId()) .build(); } @@ -85,9 +84,4 @@ public class AlarmCommentTriggerProcessor implements RuleEngineMsgNotificationRu return NotificationRuleTriggerType.ALARM_COMMENT; } - @Override - public Set getSupportedMsgTypes() { - return Set.of(DataConstants.COMMENT_CREATED, DataConstants.COMMENT_UPDATED); - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/DeviceActivityTriggerProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/DeviceActivityTriggerProcessor.java index 6e181044c7..3188eeabb1 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/DeviceActivityTriggerProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/DeviceActivityTriggerProcessor.java @@ -18,9 +18,7 @@ package org.thingsboard.server.service.notification.rule.trigger; import lombok.RequiredArgsConstructor; import org.apache.commons.collections.CollectionUtils; import org.springframework.stereotype.Service; -import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DeviceProfile; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.info.DeviceActivityNotificationInfo; @@ -28,28 +26,22 @@ import org.thingsboard.server.common.data.notification.info.RuleOriginatedNotifi import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityNotificationRuleTriggerConfig.DeviceEvent; import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.notification.trigger.RuleEngineMsgTrigger; +import org.thingsboard.server.common.msg.notification.trigger.DeviceActivityTrigger; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import java.util.Set; - @Service @RequiredArgsConstructor -public class DeviceActivityTriggerProcessor implements RuleEngineMsgNotificationRuleTriggerProcessor { +public class DeviceActivityTriggerProcessor implements NotificationRuleTriggerProcessor { private final TbDeviceProfileCache deviceProfileCache; @Override - public boolean matchesFilter(RuleEngineMsgTrigger trigger, DeviceActivityNotificationRuleTriggerConfig triggerConfig) { - if (trigger.getMsg().getOriginator().getEntityType() != EntityType.DEVICE) { - return false; - } - DeviceEvent event = trigger.getMsg().getType().equals(DataConstants.ACTIVITY_EVENT) ? DeviceEvent.ACTIVE : DeviceEvent.INACTIVE; + public boolean matchesFilter(DeviceActivityTrigger trigger, DeviceActivityNotificationRuleTriggerConfig triggerConfig) { + DeviceEvent event = trigger.isActive() ? DeviceEvent.ACTIVE : DeviceEvent.INACTIVE; if (!triggerConfig.getNotifyOn().contains(event)) { return false; } - DeviceId deviceId = (DeviceId) trigger.getMsg().getOriginator(); + DeviceId deviceId = trigger.getDeviceId(); if (CollectionUtils.isNotEmpty(triggerConfig.getDevices())) { return triggerConfig.getDevices().contains(deviceId.getId()); } else if (CollectionUtils.isNotEmpty(triggerConfig.getDeviceProfiles())) { @@ -61,15 +53,14 @@ public class DeviceActivityTriggerProcessor implements RuleEngineMsgNotification } @Override - public RuleOriginatedNotificationInfo constructNotificationInfo(RuleEngineMsgTrigger trigger) { - TbMsg msg = trigger.getMsg(); + public RuleOriginatedNotificationInfo constructNotificationInfo(DeviceActivityTrigger trigger) { return DeviceActivityNotificationInfo.builder() - .eventType(trigger.getMsg().getType().equals(DataConstants.ACTIVITY_EVENT) ? "active" : "inactive") - .deviceId(msg.getOriginator().getId()) - .deviceName(msg.getMetaData().getValue("deviceName")) - .deviceType(msg.getMetaData().getValue("deviceType")) - .deviceLabel(msg.getMetaData().getValue("deviceLabel")) - .deviceCustomerId(msg.getCustomerId()) + .eventType(trigger.isActive() ? "active" : "inactive") + .deviceId(trigger.getDeviceId().getId()) + .deviceName(trigger.getDeviceName()) + .deviceType(trigger.getDeviceType()) + .deviceLabel(trigger.getDeviceLabel()) + .deviceCustomerId(trigger.getCustomerId()) .build(); } @@ -78,9 +69,4 @@ public class DeviceActivityTriggerProcessor implements RuleEngineMsgNotification return NotificationRuleTriggerType.DEVICE_ACTIVITY; } - @Override - public Set getSupportedMsgTypes() { - return Set.of(DataConstants.ACTIVITY_EVENT, DataConstants.INACTIVITY_EVENT); - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index e33e501fbc..337d607fac 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -31,7 +31,6 @@ import org.apache.commons.lang3.tuple.Pair; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; -import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; @@ -63,6 +62,8 @@ import org.thingsboard.server.common.data.query.EntityListFilter; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; +import org.thingsboard.server.common.msg.notification.trigger.DeviceActivityTrigger; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -70,12 +71,10 @@ import org.thingsboard.server.common.stats.TbApiUsageReportClient; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.sql.query.EntityQueryRepository; -import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.util.DbTypeInfoComponent; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; @@ -158,6 +157,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService msgTypeToTriggerType; // set on init by DefaultNotificationRuleProcessor + private final AlarmInfo alarmInfo; + private final ActionType actionType; + private final User user; @Override - public NotificationRuleTriggerType getType() { - return msgTypeToTriggerType != null ? msgTypeToTriggerType.get(msg.getType()) : null; + public EntityId getOriginatorEntityId() { + return alarmInfo.getOriginator(); } @Override - public EntityId getOriginatorEntityId() { - return msg.getOriginator(); + public NotificationRuleTriggerType getType() { + return NotificationRuleTriggerType.ALARM_ASSIGNMENT; } } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/AlarmCommentTrigger.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/AlarmCommentTrigger.java new file mode 100644 index 0000000000..d0b3bdd5de --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/AlarmCommentTrigger.java @@ -0,0 +1,48 @@ +/** + * 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.common.msg.notification.trigger; + +import lombok.Builder; +import lombok.Data; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.alarm.AlarmComment; +import org.thingsboard.server.common.data.audit.ActionType; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; + +@Data +@Builder +public class AlarmCommentTrigger implements NotificationRuleTrigger { + + private final TenantId tenantId; + private final AlarmComment comment; + private final Alarm alarm; + private final ActionType actionType; + private final User user; + + @Override + public NotificationRuleTriggerType getType() { + return NotificationRuleTriggerType.ALARM_COMMENT; + } + + @Override + public EntityId getOriginatorEntityId() { + return alarm.getId(); + } + +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/ApiUsageLimitTrigger.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/ApiUsageLimitTrigger.java index 368a16712c..f21d3077ca 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/ApiUsageLimitTrigger.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/ApiUsageLimitTrigger.java @@ -36,11 +36,6 @@ public class ApiUsageLimitTrigger implements NotificationRuleTrigger { return NotificationRuleTriggerType.API_USAGE_LIMIT; } - @Override - public TenantId getTenantId() { - return tenantId; - } - @Override public EntityId getOriginatorEntityId() { return tenantId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/DeviceActivityTrigger.java b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/DeviceActivityTrigger.java new file mode 100644 index 0000000000..b426b6674b --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/notification/trigger/DeviceActivityTrigger.java @@ -0,0 +1,49 @@ +/** + * 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.common.msg.notification.trigger; + +import lombok.Builder; +import lombok.Data; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTriggerType; + +@Data +@Builder +public class DeviceActivityTrigger implements NotificationRuleTrigger { + + private final TenantId tenantId; + private final CustomerId customerId; + private final DeviceId deviceId; + private final boolean active; + + private final String deviceName; + private final String deviceType; + private final String deviceLabel; + + @Override + public EntityId getOriginatorEntityId() { + return deviceId; + } + + @Override + public NotificationRuleTriggerType getType() { + return NotificationRuleTriggerType.DEVICE_ACTIVITY; + } + +} From 184a554d08b2980258b33f2d2ee146a1337dd967 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 2 Jun 2023 15:00:03 +0300 Subject: [PATCH 4/5] More tests for notification rules; fix unstable tests --- .../AbstractNotificationApiTest.java | 38 ++++- .../notification/NotificationRuleApiTest.java | 154 +++++++++++++++++- 2 files changed, 183 insertions(+), 9 deletions(-) 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 ac3857dbdf..f45882e416 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 @@ -17,6 +17,7 @@ package org.thingsboard.server.service.notification; import com.fasterxml.jackson.core.type.TypeReference; import org.apache.commons.lang3.RandomStringUtils; +import org.junit.After; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.mock.mockito.MockBean; import org.springframework.data.util.Pair; @@ -26,6 +27,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.TenantId; import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.notification.Notification; @@ -43,6 +45,7 @@ import org.thingsboard.server.common.data.notification.settings.NotificationSett import org.thingsboard.server.common.data.notification.targets.NotificationTarget; import org.thingsboard.server.common.data.notification.targets.platform.PlatformUsersNotificationTargetConfig; import org.thingsboard.server.common.data.notification.targets.platform.UserListFilter; +import org.thingsboard.server.common.data.notification.targets.platform.UsersFilter; import org.thingsboard.server.common.data.notification.template.DeliveryMethodNotificationTemplate; import org.thingsboard.server.common.data.notification.template.EmailDeliveryMethodNotificationTemplate; import org.thingsboard.server.common.data.notification.template.NotificationTemplate; @@ -54,6 +57,10 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.dao.DaoUtil; +import org.thingsboard.server.dao.notification.NotificationRequestService; +import org.thingsboard.server.dao.notification.NotificationRuleService; +import org.thingsboard.server.dao.notification.NotificationTargetService; +import org.thingsboard.server.dao.notification.NotificationTemplateService; import java.net.URISyntaxException; import java.util.Arrays; @@ -72,21 +79,40 @@ public abstract class AbstractNotificationApiTest extends AbstractControllerTest @MockBean protected SlackService slackService; - @Autowired protected MailService mailService; + @Autowired + protected NotificationRuleService notificationRuleService; + @Autowired + protected NotificationTemplateService notificationTemplateService; + @Autowired + protected NotificationTargetService notificationTargetService; + @Autowired + protected NotificationRequestService notificationRequestService; + public static final String DEFAULT_NOTIFICATION_SUBJECT = "Just a test"; public static final NotificationType DEFAULT_NOTIFICATION_TYPE = NotificationType.GENERAL; + @After + public void afterEach() { + notificationRequestService.deleteNotificationRequestsByTenantId(TenantId.SYS_TENANT_ID); + notificationRuleService.deleteNotificationRulesByTenantId(TenantId.SYS_TENANT_ID); + notificationTemplateService.deleteNotificationTemplatesByTenantId(TenantId.SYS_TENANT_ID); + notificationTargetService.deleteNotificationTargetsByTenantId(TenantId.SYS_TENANT_ID); + } + protected NotificationTarget createNotificationTarget(UserId... usersIds) { - NotificationTarget notificationTarget = new NotificationTarget(); - notificationTarget.setTenantId(tenantId); - notificationTarget.setName("Users " + List.of(usersIds)); - PlatformUsersNotificationTargetConfig targetConfig = new PlatformUsersNotificationTargetConfig(); UserListFilter filter = new UserListFilter(); filter.setUsersIds(DaoUtil.toUUIDs(List.of(usersIds))); - targetConfig.setUsersFilter(filter); + return createNotificationTarget(filter); + } + + protected NotificationTarget createNotificationTarget(UsersFilter usersFilter) { + NotificationTarget notificationTarget = new NotificationTarget(); + notificationTarget.setName(usersFilter.toString()); + PlatformUsersNotificationTargetConfig targetConfig = new PlatformUsersNotificationTargetConfig(); + targetConfig.setUsersFilter(usersFilter); notificationTarget.setConfiguration(targetConfig); return saveNotificationTarget(notificationTarget); } diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java index 044428e270..6c4413d1e7 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java @@ -17,12 +17,14 @@ package org.thingsboard.server.service.notification; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.BooleanNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import org.junit.Before; import org.junit.Test; import org.junit.function.ThrowingRunnable; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.data.util.Pair; +import org.springframework.test.context.TestPropertySource; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; @@ -31,10 +33,15 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.UpdateMessage; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.alarm.AlarmComment; +import org.thingsboard.server.common.data.alarm.AlarmCommentType; import org.thingsboard.server.common.data.alarm.AlarmSearchStatus; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration; +import org.thingsboard.server.common.data.device.data.DefaultDeviceTransportConfiguration; +import org.thingsboard.server.common.data.device.data.DeviceData; import org.thingsboard.server.common.data.device.profile.AlarmCondition; import org.thingsboard.server.common.data.device.profile.AlarmConditionFilter; import org.thingsboard.server.common.data.device.profile.AlarmConditionFilterKey; @@ -42,6 +49,8 @@ import org.thingsboard.server.common.data.device.profile.AlarmConditionKeyType; import org.thingsboard.server.common.data.device.profile.AlarmRule; import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm; import org.thingsboard.server.common.data.device.profile.SimpleAlarmConditionSpec; +import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.notification.Notification; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.NotificationRequest; @@ -52,8 +61,11 @@ import org.thingsboard.server.common.data.notification.rule.DefaultNotificationR import org.thingsboard.server.common.data.notification.rule.EscalatedNotificationRuleRecipientsConfig; import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.notification.rule.NotificationRuleInfo; +import org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentNotificationRuleTriggerConfig; +import org.thingsboard.server.common.data.notification.rule.trigger.AlarmCommentNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.AlarmNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.AlarmNotificationRuleTriggerConfig.AlarmAction; +import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.EntitiesLimitNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.EntityActionNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.rule.trigger.NewPlatformVersionNotificationRuleTriggerConfig; @@ -68,13 +80,16 @@ import org.thingsboard.server.common.data.query.FilterPredicateValue; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.security.Authority; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.common.msg.notification.trigger.NewPlatformVersionTrigger; +import org.thingsboard.server.dao.notification.DefaultNotifications; import org.thingsboard.server.dao.notification.NotificationRequestService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.util.limits.LimitedApi; import org.thingsboard.server.dao.util.limits.RateLimitService; -import org.thingsboard.server.queue.notification.NotificationRuleProcessor; +import org.thingsboard.server.service.notification.rule.cache.DefaultNotificationRulesCache; +import org.thingsboard.server.service.state.DeviceStateService; import org.thingsboard.server.service.telemetry.AlarmSubscriptionService; import java.util.ArrayList; @@ -94,8 +109,15 @@ import static org.assertj.core.api.Assertions.offset; import static org.assertj.core.api.InstanceOfAssertFactories.type; import static org.awaitility.Awaitility.await; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentNotificationRuleTriggerConfig.Action.ASSIGNED; +import static org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentNotificationRuleTriggerConfig.Action.UNASSIGNED; +import static org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityNotificationRuleTriggerConfig.DeviceEvent.ACTIVE; +import static org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityNotificationRuleTriggerConfig.DeviceEvent.INACTIVE; @DaoSqlTest +@TestPropertySource(properties = { + "transport.http.enabled=true" +}) public class NotificationRuleApiTest extends AbstractNotificationApiTest { @SpyBean @@ -108,6 +130,12 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { private RuleChainService ruleChainService; @Autowired private NotificationRuleProcessor notificationRuleProcessor; + @Autowired + private DefaultNotifications defaultNotifications; + @Autowired + private DefaultNotificationRulesCache notificationRulesCache; + @Autowired + private DeviceStateService deviceStateService; @Before public void beforeEach() throws Exception { @@ -297,7 +325,9 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { notification = getWsClient().getLastDataUpdate().getUpdate(); assertThat(notification.getSubject()).isEqualTo("critical alarm '" + alarmType + "' is CLEARED_UNACK"); - assertThat(findNotificationRequests(EntityType.ALARM).getData()).filteredOn(NotificationRequest::isScheduled).isEmpty(); + await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> { + assertThat(findNotificationRequests(EntityType.ALARM).getData()).filteredOn(NotificationRequest::isScheduled).isEmpty(); + }); } @Test @@ -370,6 +400,122 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { }); } + @Test + public void testNotificationRuleProcessing_alarmAssignment() throws Exception { + AlarmAssignmentNotificationRuleTriggerConfig triggerConfig = AlarmAssignmentNotificationRuleTriggerConfig.builder() + .alarmTypes(Set.of("test")) + .notifyOn(Set.of(ASSIGNED, UNASSIGNED)) + .build(); + NotificationTarget target = createNotificationTarget(tenantAdminUserId); + String template = "${userEmail} ${action} alarm on ${alarmOriginatorEntityType} '${alarmOriginatorName}'. Assignee: ${assigneeEmail}"; + createNotificationRule(triggerConfig, "Test", template, target.getId()); + + Device device = createDevice("Device A", "123"); + Alarm alarm = Alarm.builder() + .tenantId(tenantId) + .originator(device.getId()) + .cleared(false) + .acknowledged(false) + .severity(AlarmSeverity.CRITICAL) + .type("test") + .startTs(System.currentTimeMillis()) + .build(); + alarm = doPost("/api/alarm", alarm, Alarm.class); + AlarmId alarmId = alarm.getId(); + + checkNotificationAfter(() -> { + doPost("/api/alarm/" + alarmId + "/assign/" + tenantAdminUserId).andExpect(status().isOk()); + }, notification -> { + assertThat(notification.getText()).isEqualTo( + TENANT_ADMIN_EMAIL + " assigned alarm on Device 'Device A'. Assignee: " + TENANT_ADMIN_EMAIL + ); + }); + + checkNotificationAfter(() -> { + doDelete("/api/alarm/" + alarmId + "/assign").andExpect(status().isOk()); + }, notification -> { + assertThat(notification.getText()).isEqualTo( + TENANT_ADMIN_EMAIL + " unassigned alarm on Device 'Device A'. Assignee: " + ); + }); + } + + @Test + public void testNotificationRuleProcessing_alarmComment() throws Exception { + AlarmCommentNotificationRuleTriggerConfig triggerConfig = AlarmCommentNotificationRuleTriggerConfig.builder() + .alarmTypes(Set.of("test")) + .onlyUserComments(true) + .notifyOnCommentUpdate(true) + .build(); + NotificationTarget target = createNotificationTarget(tenantAdminUserId); + String template = "${userEmail} ${action} comment on alarm ${alarmType}: ${comment}"; + createNotificationRule(triggerConfig, "Test", template, target.getId()); + + Device device = createDevice("Device A", "123"); + Alarm alarm = Alarm.builder() + .tenantId(tenantId) + .originator(device.getId()) + .cleared(false) + .acknowledged(false) + .severity(AlarmSeverity.CRITICAL) + .type("test") + .startTs(System.currentTimeMillis()) + .build(); + alarm = doPost("/api/alarm", alarm, Alarm.class); + AlarmId alarmId = alarm.getId(); + + AlarmComment comment = checkNotificationAfter(() -> { + return doPost("/api/alarm/" + alarmId + "/comment", + AlarmComment.builder() + .type(AlarmCommentType.OTHER) + .comment(JacksonUtil.newObjectNode() + .put("text", "this is bad")) + .build(), AlarmComment.class); + }, (notification, r) -> { + assertThat(notification.getText()).isEqualTo( + TENANT_ADMIN_EMAIL + " added comment on alarm test: this is bad" + ); + }); + + checkNotificationAfter(() -> { + ((ObjectNode) comment.getComment()).put("text", "this is very bad"); + doPost("/api/alarm/" + alarmId + "/comment", comment); + }, notification -> { + assertThat(notification.getText()).isEqualTo( + TENANT_ADMIN_EMAIL + " updated comment on alarm test: this is very bad" + ); + }); + } + + @Test + public void testNotificationRuleProcessing_deviceActivity() throws Exception { + DeviceActivityNotificationRuleTriggerConfig triggerConfig = DeviceActivityNotificationRuleTriggerConfig.builder() + .notifyOn(Set.of(ACTIVE, INACTIVE)) + .build(); + NotificationTarget target = createNotificationTarget(tenantAdminUserId); + String template = "Device ${deviceName} (${deviceLabel}) of type ${deviceType} is now ${eventType}"; + createNotificationRule(triggerConfig, "Test", template, target.getId()); + + Device device = new Device(); + device.setName("A"); + device.setLabel("Test Device A"); + device.setType("test"); + DeviceData deviceData = new DeviceData(); + deviceData.setTransportConfiguration(new DefaultDeviceTransportConfiguration()); + deviceData.setConfiguration(new DefaultDeviceConfiguration()); + device.setDeviceData(deviceData); + device = doPost("/api/device", device, Device.class); + DeviceId deviceId = device.getId(); + + checkNotificationAfter(() -> { + deviceStateService.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); + }, notification -> { + assertThat(notification.getText()).isEqualTo( + "Device A (Test Device A) of type test is now active" + ); + }); + } + @Test public void testNotificationRuleInfo() throws Exception { NotificationDeliveryMethod[] deliveryMethods = {NotificationDeliveryMethod.WEB, NotificationDeliveryMethod.EMAIL}; @@ -477,6 +623,7 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { triggerConfig.setEntityTypes(Set.of(EntityType.DEVICE)); triggerConfig.setCreated(true); NotificationRule rule = createNotificationRule(triggerConfig, "Created", "Created", createNotificationTarget(tenantAdminUserId).getId()); + notificationRulesCache.evict(tenantId); assertThat(getMyNotifications(false, 100)).size().isZero(); createDevice("Device 1", "default", "111"); @@ -487,6 +634,7 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { rule.setEnabled(false); saveNotificationRule(rule); + notificationRulesCache.evict(tenantId); createDevice("Device 2", "default", "222"); TimeUnit.SECONDS.sleep(5); @@ -494,7 +642,7 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest { rule.setEnabled(true); saveNotificationRule(rule); - TimeUnit.SECONDS.sleep(2); // for rule update event to reach rules cache + notificationRulesCache.evict(tenantId); createDevice("Device 3", "default", "333"); await().atMost(30, TimeUnit.SECONDS) From f2684e124d7afd5f5a3d7f106980d8c18cbfcd07 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 2 Jun 2023 18:20:15 +0300 Subject: [PATCH 5/5] Fix AlarmCommentControllerTest --- .../controller/AbstractNotifyEntityTest.java | 20 ++++++++-------- .../AlarmCommentControllerTest.java | 23 +++++++++++-------- 2 files changed, 23 insertions(+), 20 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java index 6c42c7604d..d5eabdff2a 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java @@ -264,7 +264,7 @@ public abstract class AbstractNotifyEntityTest extends AbstractWebTest { int cntTime = 1; testNotificationMsgToEdgeServiceTime(entityId, tenantId, actionType, cntTime); testLogEntityAction(entity, originatorId, tenantId, customerId, userId, userName, actionType, cntTime, additionalInfo); - tesPushMsgToCoreTime(cntTime); + testPushMsgToCoreTime(cntTime); Mockito.reset(tbClusterService, auditLogService); } @@ -363,13 +363,13 @@ public abstract class AbstractNotifyEntityTest extends AbstractWebTest { Mockito.any(entityId.getClass()), Mockito.any(ComponentLifecycleEvent.class)); } - private void tesPushMsgToCoreTime(int cntTime) { + private void testPushMsgToCoreTime(int cntTime) { Mockito.verify(tbClusterService, times(cntTime)).pushMsgToCore(Mockito.any(ToDeviceActorNotificationMsg.class), Mockito.isNull()); } protected void testLogEntityAction(HasName entity, EntityId originatorId, TenantId tenantId, - CustomerId customerId, UserId userId, String userName, - ActionType actionType, int cntTime, Object... additionalInfo) { + CustomerId customerId, UserId userId, String userName, + ActionType actionType, int cntTime, Object... additionalInfo) { ArgumentMatcher matcherEntityEquals = entity == null ? Objects::isNull : argument -> argument.toString().equals(entity.toString()); ArgumentMatcher matcherOriginatorId = argument -> argument.equals(originatorId); ArgumentMatcher matcherCustomerId = customerId == null ? @@ -380,10 +380,10 @@ public abstract class AbstractNotifyEntityTest extends AbstractWebTest { actionType, cntTime, extractMatcherAdditionalInfo(additionalInfo)); } - private void testLogEntityActionEntityEqClass(HasName entity, EntityId originatorId, TenantId tenantId, - CustomerId customerId, UserId userId, String userName, - ActionType actionType, int cntTime, Object... additionalInfo) { - ArgumentMatcher matcherEntityEquals = argument -> argument.getClass().equals(entity.getClass()); + protected void testLogEntityActionEntityEqClass(HasName entity, EntityId originatorId, TenantId tenantId, + CustomerId customerId, UserId userId, String userName, + ActionType actionType, int cntTime, Object... additionalInfo) { + ArgumentMatcher matcherEntityEquals = argument -> entity.getClass().isAssignableFrom(argument.getClass()); ArgumentMatcher matcherOriginatorId = argument -> argument.equals(originatorId); ArgumentMatcher matcherCustomerId = customerId == null ? argument -> argument.getClass().equals(CustomerId.class) : argument -> argument.equals(customerId); @@ -600,8 +600,8 @@ public abstract class AbstractNotifyEntityTest extends AbstractWebTest { return fieldName + " length must be equal or less than 255"; } - protected String msgErrorNoFound(String entityClassName, String assetIdStr) { - return entityClassName + " with id [" + assetIdStr + "] is not found"; + protected String msgErrorNoFound(String entityClassName, String entityIdStr) { + return entityClassName + " with id [" + entityIdStr + "] is not found"; } private String entityClassToEntityTypeName(HasName entity) { diff --git a/application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java index 287ae7cfe1..bf3e7528c9 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java @@ -46,6 +46,7 @@ import java.util.LinkedList; import java.util.List; import static org.hamcrest.Matchers.containsString; +import static org.hamcrest.Matchers.equalTo; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @Slf4j @@ -104,7 +105,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { AlarmComment createdComment = createAlarmComment(alarm.getId()); - testLogEntityAction(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.ADDED_COMMENT, 1, createdComment); + testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.ADDED_COMMENT, 1, createdComment); } @Test @@ -116,7 +117,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { AlarmComment createdComment = createAlarmComment(alarm.getId()); Assert.assertEquals(AlarmCommentType.OTHER, createdComment.getType()); - testLogEntityAction(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.ADDED_COMMENT, 1, createdComment); + testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.ADDED_COMMENT, 1, createdComment); } @Test @@ -135,7 +136,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { Assert.assertEquals("true", updatedAlarmComment.getComment().get("edited").asText()); Assert.assertNotNull(updatedAlarmComment.getComment().get("editedOn")); - testLogEntityAction(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.UPDATED_COMMENT, 1, savedComment); + testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.UPDATED_COMMENT, 1, savedComment); } @Test @@ -154,7 +155,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { Assert.assertEquals("true", updatedAlarmComment.getComment().get("edited").asText()); Assert.assertNotNull(updatedAlarmComment.getComment().get("editedOn")); - testLogEntityAction(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.UPDATED_COMMENT, 1, updatedAlarmComment); + testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.UPDATED_COMMENT, 1, updatedAlarmComment); } @Test @@ -169,8 +170,8 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { savedComment.setComment(newComment); doPost("/api/alarm/" + alarm.getId() + "/comment", savedComment) - .andExpect(status().isForbidden()) - .andExpect(statusReason(containsString(msgErrorPermission))); + .andExpect(status().isNotFound()) + .andExpect(statusReason(equalTo(msgErrorNoFound("Alarm", alarm.getId().toString())))); testNotifyEntityNever(alarm.getId(), savedComment); } @@ -209,7 +210,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { .comment(JacksonUtil.newObjectNode().put("text", String.format("User %s deleted his comment", CUSTOMER_USER_EMAIL))) .build(); - testLogEntityAction(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.DELETED_COMMENT, 1, expectedAlarmComment); + testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.DELETED_COMMENT, 1, expectedAlarmComment); } @Test @@ -228,7 +229,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { .comment(JacksonUtil.newObjectNode().put("text", String.format("User %s deleted his comment", TENANT_ADMIN_EMAIL))) .build(); - testLogEntityAction(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.DELETED_COMMENT, 1, expectedAlarmComment); + testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.DELETED_COMMENT, 1, expectedAlarmComment); } @Test @@ -355,16 +356,18 @@ public class AlarmCommentControllerTest extends AbstractControllerTest { Assert.assertTrue("Created alarm doesn't match the found one!", equals); } - private AlarmComment createAlarmComment(AlarmId alarmId, String text) { + private AlarmComment createAlarmComment(AlarmId alarmId, String text) { AlarmComment alarmComment = AlarmComment.builder() .comment(JacksonUtil.newObjectNode().set("text", new TextNode(text))) .build(); return saveAlarmComment(alarmId, alarmComment); } - private AlarmComment createAlarmComment(AlarmId alarmId) { + + private AlarmComment createAlarmComment(AlarmId alarmId) { return createAlarmComment(alarmId, "Please take a look"); } + private AlarmComment saveAlarmComment(AlarmId alarmId, AlarmComment alarmComment) { alarmComment = doPost("/api/alarm/" + alarmId + "/comment", alarmComment, AlarmComment.class); Assert.assertNotNull(alarmComment);