diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index 019b6eb5e8..cf471db8ee 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -32,13 +32,16 @@ import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.springframework.web.socket.CloseStatus; import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.ComparisonTsValue; +import org.thingsboard.server.common.data.query.OriginatorAlarmFilter; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityKey; @@ -52,6 +55,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.ws.WebSocketService; import org.thingsboard.server.service.ws.WebSocketSessionRef; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggHistoryCmd; @@ -60,6 +64,8 @@ import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggTimeSeriesCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataUpdate; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdate; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataUpdate; @@ -68,6 +74,7 @@ import org.thingsboard.server.service.ws.telemetry.cmd.v2.GetTsCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.LatestValueCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.TimeSeriesCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.UnsubscribeCmd; +import org.thingsboard.server.service.ws.telemetry.sub.AlarmSubscriptionUpdate; import java.util.ArrayList; import java.util.Arrays; @@ -76,6 +83,7 @@ import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; @@ -139,6 +147,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private int maxAlarmQueriesPerRefreshInterval; @Value("${ui.dashboard.max_datapoints_limit:50000}") private int maxDatapointLimit; + @Value("${server.ws.alarms_per_alarm_status_subscription_cache_size:10}") + private int alarmsPerAlarmStatusSubscriptionCacheSize; private ExecutorService wsCallBackExecutor; private boolean tsInSqlDB; @@ -434,6 +444,76 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } + @Override + public void handleCmd(WebSocketSessionRef sessionRef, AlarmStatusCmd cmd) { + log.debug("[{}] Handling alarm status subscription cmd (cmdId: {})", sessionRef.getSessionId(), cmd.getCmdId()); + SecurityUser securityCtx = sessionRef.getSecurityCtx(); + + TbAlarmStatusSubscription subscription = TbAlarmStatusSubscription.builder() + .serviceId(serviceInfoProvider.getServiceId()) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(cmd.getCmdId()) + .tenantId(securityCtx.getTenantId()) + .entityId(cmd.getOriginatorId()) + .typeList(cmd.getTypeList()) + .severityList(cmd.getSeverityList()) + .updateProcessor(this::handleAlarmStatusSubscriptionUpdate) + .build(); + localSubscriptionService.addSubscription(subscription, sessionRef); + + fetchActiveAlarms(subscription); + sendUpdate(sessionRef.getSessionId(), subscription.createUpdate()); + } + + private void fetchActiveAlarms(TbAlarmStatusSubscription subscription) { + log.trace("[{}, subId: {}] Fetching active alarms from DB", subscription.getSessionId(), subscription.getSubscriptionId()); + OriginatorAlarmFilter originatorAlarmFilter = new OriginatorAlarmFilter(subscription.getEntityId(), subscription.getTypeList(), subscription.getSeverityList()); + List alarmIds = alarmService.findActiveOriginatorAlarms(subscription.getTenantId(), originatorAlarmFilter, new PageLink(alarmsPerAlarmStatusSubscriptionCacheSize)).getData(); + + subscription.getAlarmIds().addAll(alarmIds); + subscription.setExceededLimit(alarmIds.size() == alarmsPerAlarmStatusSubscriptionCacheSize); + } + + private void sendUpdate(String sessionId, CmdUpdate update) { + log.trace("[{}, cmdId: {}] Sending WS update: {}", sessionId, update.getCmdId(), update); + wsService.sendUpdate(sessionId, update); + } + + private void handleAlarmStatusSubscriptionUpdate(TbSubscription sub, AlarmSubscriptionUpdate subscriptionUpdate) { + TbAlarmStatusSubscription subscription = (TbAlarmStatusSubscription) sub; + try { + AlarmInfo alarm = subscriptionUpdate.getAlarm(); + Set alarmsIds = subscription.getAlarmIds(); + if (alarmsIds.contains(alarm.getId().getId())) { + if (!alarmMatchesSubscription(alarm, subscription) || subscriptionUpdate.isAlarmDeleted()) { + alarmsIds.remove(alarm.getId().getId()); + if (alarmsIds.size() == 0) { + if (subscription.isExceededLimit()) { + fetchActiveAlarms(subscription); + if (alarmsIds.size() == 0) { + sendUpdate(subscription.getSessionId(), subscription.createUpdate()); + } + } else { + sendUpdate(subscription.getSessionId(), subscription.createUpdate()); + } + } + } + } else if (alarmMatchesSubscription(alarm, subscription) && (alarmsIds.size() < alarmsPerAlarmStatusSubscriptionCacheSize)) { + alarmsIds.add(alarm.getId().getId()); + if (alarmsIds.size() == 1) { + sendUpdate(subscription.getSessionId(), subscription.createUpdate()); + } + } + } catch (Exception e) { + log.error("[{}, subId: {}] Failed to handle update for alarm status subscription: {}", subscription.getSessionId(), subscription.getSubscriptionId(), subscriptionUpdate, e); + } + } + + private boolean alarmMatchesSubscription(AlarmInfo alarm, TbAlarmStatusSubscription subscription) { + return !alarm.isCleared() && (subscription.getTypeList() == null || subscription.getTypeList().contains(alarm.getType())) && + (subscription.getSeverityList() == null || subscription.getSeverityList().contains(alarm.getSeverity())); + } + private boolean validate(TbAbstractSubCtx finalCtx) { if (finalCtx.isStopped()) { log.warn("[{}][{}][{}] Received validation task for already stopped context.", finalCtx.getTenantId(), finalCtx.getSessionId(), finalCtx.getCmdId()); @@ -527,7 +607,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc return ctx; } - private TbAlarmCountSubCtx createSubCtx(WebSocketSessionRef sessionRef, AlarmCountCmd cmd) { + private TbAlarmCountSubCtx createSubCtx(WebSocketSessionRef sessionRef, AlarmCountCmd cmd) { Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new ConcurrentHashMap<>()); TbAlarmCountSubCtx ctx = new TbAlarmCountSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, alarmService, sessionRef, cmd.getCmdId()); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java index 0a47f5f375..a8f8c51f03 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java @@ -414,7 +414,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer private void onAlarmUpdate(UUID entityId, AlarmSubscriptionUpdate update, TbCallback callback) { processSubscriptionData(entityId, - sub -> TbSubscriptionType.ALARMS.equals(sub.getType()), + sub -> TbSubscriptionType.ALARMS.equals(sub.getType()) || TbSubscriptionType.ALARM_STATUS.equals(sub.getType()), update, callback); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmStatusSubscription.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmStatusSubscription.java new file mode 100644 index 0000000000..6b16383378 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmStatusSubscription.java @@ -0,0 +1,71 @@ +/** + * Copyright © 2016-2024 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.service.subscription; + +import lombok.Builder; +import lombok.Getter; +import lombok.Setter; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusUpdate; +import org.thingsboard.server.service.ws.telemetry.sub.AlarmSubscriptionUpdate; + +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.UUID; +import java.util.function.BiConsumer; + + +public class TbAlarmStatusSubscription extends TbSubscription { + + @Getter + private final Set alarmIds = new HashSet<>(); + @Getter + @Setter + private boolean exceededLimit; + @Getter + private final List typeList; + @Getter + private final List severityList; + + @Builder + public TbAlarmStatusSubscription(String serviceId, String sessionId, int subscriptionId, TenantId tenantId, EntityId entityId, + BiConsumer, AlarmSubscriptionUpdate> updateProcessor, + List typeList, List severityList) { + super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.ALARM_STATUS, updateProcessor); + this.typeList = typeList; + this.severityList = severityList; + } + + @Override + public boolean equals(Object o) { + return super.equals(o); + } + + @Override + public int hashCode() { + return super.hashCode(); + } + + public AlarmStatusUpdate createUpdate() { + return AlarmStatusUpdate.builder() + .cmdId(getSubscriptionId()) + .present(alarmIds.size() > 0) + .build(); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java index 865602f7a7..1d8608b31d 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.subscription; import org.thingsboard.server.service.ws.WebSocketSessionRef; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataCmd; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.UnsubscribeCmd; @@ -32,6 +33,8 @@ public interface TbEntityDataSubscriptionService { void handleCmd(WebSocketSessionRef sessionId, AlarmCountCmd cmd); + void handleCmd(WebSocketSessionRef session, AlarmStatusCmd cmd); + void cancelSubscription(String sessionId, UnsubscribeCmd subscriptionId); void cancelAllSessionSubscriptions(String sessionId); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityLocalSubsInfo.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityLocalSubsInfo.java index ee20843538..7147bfe4e8 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityLocalSubsInfo.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityLocalSubsInfo.java @@ -74,6 +74,7 @@ public class TbEntityLocalSubsInfo { stateChanged = true; } break; + case ALARM_STATUS: case ALARMS: if (!newState.alarms) { newState.alarms = true; @@ -168,6 +169,7 @@ public class TbEntityLocalSubsInfo { case NOTIFICATIONS_COUNT: state.notifications = false; break; + case ALARM_STATUS: case ALARMS: state.alarms = false; break; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionType.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionType.java index e69c06df01..a2c45a4adc 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionType.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionType.java @@ -16,5 +16,5 @@ package org.thingsboard.server.service.subscription; public enum TbSubscriptionType { - TIMESERIES, ATTRIBUTES, ALARMS, NOTIFICATIONS, NOTIFICATIONS_COUNT + TIMESERIES, ATTRIBUTES, ALARMS, ALARM_STATUS, NOTIFICATIONS, NOTIFICATIONS_COUNT } diff --git a/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java index 8d96d14b49..720a4ea8ac 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java @@ -77,6 +77,7 @@ import org.thingsboard.server.service.ws.telemetry.cmd.v1.TelemetryPluginCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v1.TimeseriesSubscriptionCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataCmd; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdate; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataCmd; @@ -168,6 +169,7 @@ public class DefaultWebSocketService implements WebSocketService { cmdsHandlers.put(WsCmdType.ALARM_DATA, newCmdHandler(this::handleWsAlarmDataCmd)); cmdsHandlers.put(WsCmdType.ENTITY_COUNT, newCmdHandler(this::handleWsEntityCountCmd)); cmdsHandlers.put(WsCmdType.ALARM_COUNT, newCmdHandler(this::handleWsAlarmCountCmd)); + cmdsHandlers.put(WsCmdType.ALARM_STATUS, newCmdHandler(this::handleWsAlarmsStatusCmd)); cmdsHandlers.put(WsCmdType.ENTITY_DATA_UNSUBSCRIBE, newCmdHandler(this::handleWsDataUnsubscribeCmd)); cmdsHandlers.put(WsCmdType.ALARM_DATA_UNSUBSCRIBE, newCmdHandler(this::handleWsDataUnsubscribeCmd)); cmdsHandlers.put(WsCmdType.ENTITY_COUNT_UNSUBSCRIBE, newCmdHandler(this::handleWsDataUnsubscribeCmd)); @@ -261,6 +263,12 @@ public class DefaultWebSocketService implements WebSocketService { } } + private void handleWsAlarmsStatusCmd(WebSocketSessionRef sessionRef, AlarmStatusCmd cmd) { + if (validateCmd(sessionRef, cmd)) { + entityDataSubService.handleCmd(sessionRef, cmd); + } + } + @Override public void sendUpdate(String sessionId, int cmdId, TelemetrySubscriptionUpdate update) { // We substitute the subscriptionId with cmdId for old-style subscriptions. diff --git a/application/src/main/java/org/thingsboard/server/service/ws/WsCmdType.java b/application/src/main/java/org/thingsboard/server/service/ws/WsCmdType.java index 54d43e7817..f591f777d3 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/WsCmdType.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/WsCmdType.java @@ -35,5 +35,6 @@ public enum WsCmdType { ALARM_COUNT_UNSUBSCRIBE, ENTITY_DATA_UNSUBSCRIBE, ENTITY_COUNT_UNSUBSCRIBE, - NOTIFICATIONS_UNSUBSCRIBE + NOTIFICATIONS_UNSUBSCRIBE, + ALARM_STATUS } diff --git a/application/src/main/java/org/thingsboard/server/service/ws/WsCommandsWrapper.java b/application/src/main/java/org/thingsboard/server/service/ws/WsCommandsWrapper.java index 694f3a7920..b79d6b2a13 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/WsCommandsWrapper.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/WsCommandsWrapper.java @@ -33,6 +33,7 @@ import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountUnsubscribeCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataUnsubscribeCmd; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountUnsubscribeCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataCmd; @@ -56,6 +57,7 @@ public class WsCommandsWrapper { @Type(name = "ENTITY_COUNT", value = EntityCountCmd.class), @Type(name = "ALARM_DATA", value = AlarmDataCmd.class), @Type(name = "ALARM_COUNT", value = AlarmCountCmd.class), + @Type(name = "ALARM_STATUS", value = AlarmStatusCmd.class), @Type(name = "NOTIFICATIONS", value = NotificationsSubCmd.class), @Type(name = "NOTIFICATIONS_COUNT", value = NotificationsCountSubCmd.class), @Type(name = "MARK_NOTIFICATIONS_AS_READ", value = MarkNotificationsAsReadCmd.class), diff --git a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusCmd.java b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusCmd.java new file mode 100644 index 0000000000..8daab4b40d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusCmd.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2024 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.service.ws.telemetry.cmd.v2; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.service.ws.WsCmd; +import org.thingsboard.server.service.ws.WsCmdType; + +import java.util.List; + +@Data +@AllArgsConstructor +@NoArgsConstructor +public class AlarmStatusCmd implements WsCmd { + + private int cmdId; + private EntityId originatorId; + private List typeList; + private List severityList; + + @Override + public WsCmdType getType() { + return WsCmdType.ALARM_STATUS; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusUpdate.java b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusUpdate.java new file mode 100644 index 0000000000..daf7af74a3 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusUpdate.java @@ -0,0 +1,54 @@ +/** + * Copyright © 2016-2024 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.service.ws.telemetry.cmd.v2; + +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Builder; +import lombok.Getter; +import lombok.ToString; +import org.thingsboard.server.service.subscription.SubscriptionErrorCode; + +@ToString +@Getter +public class AlarmStatusUpdate extends CmdUpdate { + + @Getter + private boolean present; + + public AlarmStatusUpdate(int cmdId, boolean present) { + super(cmdId, SubscriptionErrorCode.NO_ERROR.getCode(), null); + this.present = present; + } + + public AlarmStatusUpdate(int cmdId, int errorCode, String errorMsg) { + super(cmdId, errorCode, errorMsg); + } + + @Builder + public AlarmStatusUpdate(@JsonProperty("cmdId") int cmdId, + @JsonProperty("present") boolean present, + @JsonProperty("errorCode") int errorCode, + @JsonProperty("errorMsg") String errorMsg) { + super(cmdId, errorCode, errorMsg); + this.present = present; + } + + @Override + public CmdUpdateType getCmdUpdateType() { + return CmdUpdateType.ALARM_STATUS; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/CmdUpdateType.java b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/CmdUpdateType.java index c6eb494946..2366b9b991 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/CmdUpdateType.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/CmdUpdateType.java @@ -19,6 +19,7 @@ public enum CmdUpdateType { ENTITY_DATA, ALARM_DATA, ALARM_COUNT_DATA, + ALARM_STATUS, COUNT_DATA, NOTIFICATIONS, NOTIFICATIONS_COUNT diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 0c4648c7e6..2c82d544d4 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -87,6 +87,8 @@ server: subscriptions_per_tenant: "${TB_SERVER_WS_SUBSCRIPTIONS_PER_TENANT_RATE_LIMIT:}" # Per-user rate limit for WS subscriptions subscriptions_per_user: "${TB_SERVER_WS_SUBSCRIPTIONS_PER_USER_RATE_LIMIT:}" + # Maximum number of active originator alarm ids being saved in cache for single alarm status subscription. For example, no more than 10 alarm ids on the alarm widget + alarms_per_alarm_status_subscription_cache_size: "${TB_ALARMS_PER_ALARM_STATUS_SUBSCRIPTION_CACHE_SIZE:10}" rest: server_side_rpc: # Minimum value of the server-side RPC timeout. May override value provided in the REST API call. diff --git a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java index be01ee4fe8..2b988d0058 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java @@ -26,6 +26,8 @@ import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.test.context.TestPropertySource; +import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.alarm.Alarm; @@ -59,10 +61,13 @@ import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountUpdate; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusUpdate; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountUpdate; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataUpdate; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; @@ -75,6 +80,9 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers. @Slf4j @DaoSqlTest +@TestPropertySource(properties = { + "server.ws.alarms_per_alarm_status_subscription_cache_size=5" +}) public class WebsocketApiTest extends AbstractControllerTest { @Autowired private TelemetrySubscriptionService tsService; @@ -314,6 +322,119 @@ public class WebsocketApiTest extends AbstractControllerTest { Assert.assertEquals(1, update.getCount()); } + @Test + public void testAlarmStatusWsCmd() throws Exception { + loginTenantAdmin(); + + AlarmStatusCmd cmd = new AlarmStatusCmd(1, device.getId(), List.of("TEST ALARM", "TEST ALARM 2"), List.of(AlarmSeverity.WARNING)); + + getWsClient().send(cmd); + + AlarmStatusUpdate update = JacksonUtil.fromString(getWsClient().waitForReply(), AlarmStatusUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + Assert.assertFalse(update.isPresent()); + + //create alarm + getWsClient().registerWaitForUpdate(); + + Alarm alarm = new Alarm(); + alarm.setOriginator(device.getId()); + alarm.setType("TEST ALARM"); + alarm.setSeverity(AlarmSeverity.WARNING); + + alarm = doPost("/api/alarm", alarm, Alarm.class); + + AlarmStatusUpdate alarmStatusUpdate = JacksonUtil.fromString(getWsClient().waitForUpdate(), AlarmStatusUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + Assert.assertTrue(alarmStatusUpdate.isPresent()); + + //clear alarm + getWsClient().registerWaitForUpdate(); + + String alarmId = alarm.getId().getId().toString(); + Alarm clearedAlarm = doPost("/api/alarm/" + alarmId + "/clear", Alarm.class); + Assert.assertNotNull(clearedAlarm); + Assert.assertTrue(clearedAlarm.isCleared()); + + AlarmStatusUpdate alarmStatusUpdate2 = JacksonUtil.fromString(getWsClient().waitForUpdate(), AlarmStatusUpdate.class); + Assert.assertEquals(1, alarmStatusUpdate2.getCmdId()); + Assert.assertFalse(alarmStatusUpdate2.isPresent()); + + // add second type alarm + getWsClient().registerWaitForUpdate(); + + Alarm alarm2 = new Alarm(); + alarm2.setOriginator(device.getId()); + alarm2.setType("TEST ALARM 2"); + alarm2.setSeverity(AlarmSeverity.WARNING); + + doPost("/api/alarm", alarm2, Alarm.class); + + AlarmStatusUpdate alarmStatusUpdate3 = JacksonUtil.fromString(getWsClient().waitForReply(), AlarmStatusUpdate.class); + Assert.assertEquals(1, alarmStatusUpdate3.getCmdId()); + Assert.assertTrue(alarmStatusUpdate3.isPresent()); + + //change severity + alarm2.setSeverity(AlarmSeverity.MAJOR); + Alarm updatedAlarm = doPost("/api/alarm", alarm2, Alarm.class); + Assert.assertNotNull(updatedAlarm); + Assert.assertEquals(AlarmSeverity.MAJOR, updatedAlarm.getSeverity()); + + AlarmStatusUpdate alarmStatusUpdate4 = JacksonUtil.fromString(getWsClient().waitForReply(), AlarmStatusUpdate.class); + Assert.assertEquals(1, alarmStatusUpdate4.getCmdId()); + Assert.assertFalse(alarmStatusUpdate4.isPresent()); + + //subscribe for critical alarms + AlarmStatusCmd cmd3 = new AlarmStatusCmd(2, device.getId(), List.of("TEST ALARM"), List.of(AlarmSeverity.CRITICAL)); + + getWsClient().send(cmd3); + + AlarmStatusUpdate alarmStatusUpdate5 = JacksonUtil.fromString(getWsClient().waitForReply(), AlarmStatusUpdate.class); + Assert.assertEquals(2, alarmStatusUpdate5.getCmdId()); + Assert.assertFalse(alarmStatusUpdate5.isPresent()); + } + + @Test + public void testAlarmStatusWsCmdWithMaxAlarmsCacheSize() throws Exception { + loginTenantAdmin(); + + //create 5+1 alarms + List alarms = new ArrayList<>(); + for (int i = 0; i < 6; i++) { + Alarm alarm = new Alarm(); + alarm.setOriginator(device.getId()); + alarm.setType(RandomStringUtils.randomAlphabetic(10)); + alarm.setSeverity(AlarmSeverity.CRITICAL); + alarm = doPost("/api/alarm", alarm, Alarm.class); + alarms.add(alarm); + } + + AlarmStatusCmd cmd = new AlarmStatusCmd(1, device.getId(), null, List.of(AlarmSeverity.CRITICAL)); + + getWsClient().send(cmd); + + AlarmStatusUpdate update = JacksonUtil.fromString(getWsClient().waitForReply(), AlarmStatusUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + Assert.assertTrue(update.isPresent()); + + getWsClient().registerWaitForUpdate(); + //clear first 5 alarms + for (int i = 0; i < 5; i++) { + String alarmId = alarms.get(i).getId().getId().toString(); + doPost("/api/alarm/" + alarmId + "/clear", Alarm.class); + } + AlarmStatusUpdate alarmStatusUpdate = JacksonUtil.fromString(getWsClient().waitForUpdate(), AlarmStatusUpdate.class); + Assert.assertNull(alarmStatusUpdate); + + //clear 6-th alarm should send update + String alarmId6 = alarms.get(5).getId().getId().toString(); + doPost("/api/alarm/" + alarmId6 + "/clear", Alarm.class); + + AlarmStatusUpdate alarmStatusUpdate2 = JacksonUtil.fromString(getWsClient().waitForUpdate(), AlarmStatusUpdate.class); + Assert.assertEquals(1, alarmStatusUpdate2.getCmdId()); + Assert.assertFalse(alarmStatusUpdate2.isPresent()); + } + @Test public void testEntityDataLatestWidgetFlow() throws Exception { List keys = List.of(new EntityKey(EntityKeyType.TIME_SERIES, "temperature")); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java index 256d10465d..e3499ee97a 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java @@ -38,6 +38,7 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.query.AlarmCountQuery; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; +import org.thingsboard.server.common.data.query.OriginatorAlarmFilter; import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.dao.entity.EntityDaoService; @@ -119,4 +120,6 @@ public interface AlarmService extends EntityDaoService { PageData findAlarmTypesByTenantId(TenantId tenantId, PageLink pageLink); + PageData findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, PageLink pageLink); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/OriginatorAlarmFilter.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/OriginatorAlarmFilter.java new file mode 100644 index 0000000000..b62b285797 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/OriginatorAlarmFilter.java @@ -0,0 +1,33 @@ +/** + * Copyright © 2016-2024 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.data.query; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.id.EntityId; + +import java.util.List; + +@NoArgsConstructor +@AllArgsConstructor +@Getter +public class OriginatorAlarmFilter { + private EntityId originatorId; + private List typeList; + private List severityList; +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java index 07b670f835..59b882604a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java @@ -38,6 +38,7 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.query.AlarmCountQuery; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; +import org.thingsboard.server.common.data.query.OriginatorAlarmFilter; import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.dao.Dao; @@ -110,4 +111,7 @@ public interface AlarmDao extends Dao { PageData findTenantAlarmTypes(UUID tenantId, PageLink pageLink); boolean removeAlarmTypesIfNoAlarmsPresent(UUID tenantId, Set types); + + PageData findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, PageLink pageLink); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java index 7e9cabf5b4..7026de2409 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java @@ -54,6 +54,7 @@ import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.query.AlarmCountQuery; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; +import org.thingsboard.server.common.data.query.OriginatorAlarmFilter; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.EntitySearchDirection; @@ -365,6 +366,12 @@ public class BaseAlarmService extends AbstractCachedEntityService findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, PageLink pageLink) { + log.trace("Executing findActiveOriginatorAlarms, tenantId [{}], originatorAlarmFilter [{}]", tenantId, originatorAlarmFilter); + return alarmDao.findActiveOriginatorAlarms(tenantId, originatorAlarmFilter, pageLink); + } + private Alarm merge(Alarm existing, Alarm alarm) { if (alarm.getStartTs() > existing.getEndTs()) { existing.setEndTs(alarm.getStartTs()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java index 65f46ecc64..1382deb84b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java @@ -101,10 +101,8 @@ public interface AlarmRepository extends JpaRepository { "AND ea.entityType = :affectedEntityType " + "AND (:startTime IS NULL OR (a.createdTime >= :startTime AND ea.createdTime >= :startTime)) " + "AND (:endTime IS NULL OR (a.createdTime <= :endTime AND ea.createdTime <= :endTime)) " + - "AND ((:#{#alarmTypes == null} = true) OR a.type IN (:alarmTypes)) " + //HHH-15968 - "AND ((:#{#alarmSeverities == null} = true) OR a.severity IN (:alarmSeverities)) " + //HHH-15968 -// "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + -// "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + "AND ((:clearFilterEnabled) = FALSE OR a.cleared = :clearFilter) " + "AND ((:ackFilterEnabled) = FALSE OR a.acknowledged = :ackFilter) " + "AND (:assigneeId IS NULL OR a.assigneeId = :assigneeId) " + @@ -122,10 +120,8 @@ public interface AlarmRepository extends JpaRepository { "AND ea.entityType = :affectedEntityType " + "AND (:startTime IS NULL OR (a.createdTime >= :startTime AND ea.createdTime >= :startTime)) " + "AND (:endTime IS NULL OR (a.createdTime <= :endTime AND ea.createdTime <= :endTime)) " + - "AND ((:#{#alarmTypes == null} = true) OR a.type IN (:alarmTypes)) " + //HHH-15968 - "AND ((:#{#alarmSeverities == null} = true) OR a.severity IN (:alarmSeverities)) " + //HHH-15968 -// "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + -// "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + "AND ((:clearFilterEnabled) = FALSE OR a.cleared = :clearFilter) " + "AND ((:ackFilterEnabled) = FALSE OR a.acknowledged = :ackFilter) " + "AND (:assigneeId IS NULL OR a.assigneeId = :assigneeId) " + @@ -186,10 +182,8 @@ public interface AlarmRepository extends JpaRepository { "WHERE a.tenantId = :tenantId " + "AND (:startTime IS NULL OR a.createdTime >= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND ((:#{#alarmTypes == null} = true) OR a.type IN (:alarmTypes)) " + //HHH-15968 - "AND ((:#{#alarmSeverities == null} = true) OR a.severity IN (:alarmSeverities)) " + //HHH-15968 -// "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + -// "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + "AND ((:clearFilterEnabled) = FALSE OR a.cleared = :clearFilter) " + "AND ((:ackFilterEnabled) = FALSE OR a.acknowledged = :ackFilter) " + "AND (:assigneeId IS NULL OR a.assigneeId = :assigneeId) " + @@ -202,10 +196,8 @@ public interface AlarmRepository extends JpaRepository { "WHERE a.tenantId = :tenantId " + "AND (:startTime IS NULL OR a.createdTime >= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND ((:#{#alarmTypes == null} = true) OR a.type IN (:alarmTypes)) " + //HHH-15968 - "AND ((:#{#alarmSeverities == null} = true) OR a.severity IN (:alarmSeverities)) " + //HHH-15968 -// "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + -// "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + "AND ((:clearFilterEnabled) = FALSE OR a.cleared = :clearFilter) " + "AND ((:ackFilterEnabled) = FALSE OR a.acknowledged = :ackFilter) " + "AND (:assigneeId IS NULL OR a.assigneeId = :assigneeId) " + @@ -266,10 +258,8 @@ public interface AlarmRepository extends JpaRepository { "WHERE a.tenantId = :tenantId AND a.customerId = :customerId " + "AND (:startTime IS NULL OR a.createdTime >= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND ((:#{#alarmTypes == null} = true) OR a.type IN (:alarmTypes)) " + //HHH-15968 - "AND ((:#{#alarmSeverities == null} = true) OR a.severity IN (:alarmSeverities)) " + //HHH-15968 -// "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + -// "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + "AND ((:clearFilterEnabled) = FALSE OR a.cleared = :clearFilter) " + "AND ((:ackFilterEnabled) = FALSE OR a.acknowledged = :ackFilter) " + "AND (:assigneeId IS NULL OR a.assigneeId = :assigneeId) " + @@ -283,10 +273,8 @@ public interface AlarmRepository extends JpaRepository { "WHERE a.tenantId = :tenantId AND a.customerId = :customerId " + "AND (:startTime IS NULL OR a.createdTime >= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND ((:#{#alarmTypes == null} = true) OR a.type IN (:alarmTypes)) " + //HHH-15968 - "AND ((:#{#alarmSeverities == null} = true) OR a.severity IN (:alarmSeverities)) " + //HHH-15968 -// "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + -// "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + "AND ((:clearFilterEnabled) = FALSE OR a.cleared = :clearFilter) " + "AND ((:ackFilterEnabled) = FALSE OR a.acknowledged = :ackFilter) " + "AND (:assigneeId IS NULL OR a.assigneeId = :assigneeId) " + @@ -404,4 +392,25 @@ public interface AlarmRepository extends JpaRepository { @Query(value = "DELETE FROM alarm_types AS at WHERE NOT EXISTS (SELECT 1 FROM alarm AS a WHERE a.tenant_id = at.tenant_id AND a.type = at.type) AND at.tenant_id = :tenantId AND at.type IN (:types)", nativeQuery = true) int deleteTypeIfNoAlarmsExist(@Param("tenantId") UUID tenantId, @Param("types") Set types); + @Query(value = "SELECT a.id " + + "FROM AlarmEntity a " + + "WHERE a.tenantId = :tenantId " + + "AND a.originatorId = :originatorId " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND (a.cleared = false)", + countQuery = "" + + "SELECT count(a) " + + "FROM AlarmEntity a " + + "WHERE a.tenantId = :tenantId " + + "AND a.originatorId = :originatorId " + + "AND ((:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " + + "AND ((:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " + + "AND (a.cleared = false)") + Page findActiveOriginatorAlarms(@Param("tenantId") UUID tenantId, + @Param("originatorId") UUID originatorId, + @Param("alarmTypes") List alarmTypes, + @Param("alarmSeverities") List alarmSeverities, + Pageable pageable); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java index b4c1bf7913..dbf7fde100 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java @@ -54,6 +54,7 @@ import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.query.AlarmCountQuery; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; +import org.thingsboard.server.common.data.query.OriginatorAlarmFilter; import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.alarm.AlarmDao; @@ -434,6 +435,12 @@ public class JpaAlarmDao extends JpaAbstractDao implements A return alarmRepository.deleteTypeIfNoAlarmsExist(tenantId, types) > 0; } + @Override + public PageData findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, PageLink pageLink) { + return DaoUtil.pageToPageData(alarmRepository.findActiveOriginatorAlarms(tenantId.getId(), originatorAlarmFilter.getOriginatorId().getId(), + originatorAlarmFilter.getTypeList(), originatorAlarmFilter.getSeverityList(), toPageable(pageLink, false))); + } + private static String getPropagationTypes(AlarmPropagationInfo ap) { String propagateRelationTypes; if (!CollectionUtils.isEmpty(ap.getPropagateRelationTypes())) {