Browse Source

Merge branch 'alarmStatusWsCmd' of github.com:dashevchenko/thingsboard

pull/12094/head
Andrii Shvaika 2 years ago
parent
commit
c498b26485
  1. 78
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  3. 67
      application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmStatusSubscription.java
  4. 3
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityLocalSubsInfo.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionType.java
  7. 8
      application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java
  8. 3
      application/src/main/java/org/thingsboard/server/service/ws/WsCmdType.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/ws/WsCommandsWrapper.java
  10. 42
      application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusCmd.java
  11. 54
      application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/AlarmStatusUpdate.java
  12. 1
      application/src/main/java/org/thingsboard/server/service/ws/telemetry/cmd/v2/CmdUpdateType.java
  13. 2
      application/src/main/resources/thingsboard.yml
  14. 126
      application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java
  15. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java
  16. 33
      common/data/src/main/java/org/thingsboard/server/common/data/query/OriginatorAlarmFilter.java
  17. 4
      dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java
  18. 7
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
  19. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java
  20. 8
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java

78
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -32,6 +32,7 @@ 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;
@ -39,6 +40,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData;
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 +54,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 +63,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 +73,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 +82,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 +146,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 +443,75 @@ 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<UUID> alarmIds = alarmService.findActiveOriginatorAlarms(subscription.getTenantId(), originatorAlarmFilter, alarmsPerAlarmStatusSubscriptionCacheSize);
subscription.getAlarmIds().addAll(alarmIds);
subscription.setCacheFull(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<AlarmSubscriptionUpdate> sub, AlarmSubscriptionUpdate subscriptionUpdate) {
TbAlarmStatusSubscription subscription = (TbAlarmStatusSubscription) sub;
try {
AlarmInfo alarm = subscriptionUpdate.getAlarm();
Set<UUID> alarmsIds = subscription.getAlarmIds();
if (alarmsIds.contains(alarm.getId().getId())) {
if (!subscription.matches(alarm) || subscriptionUpdate.isAlarmDeleted()) {
alarmsIds.remove(alarm.getId().getId());
if (alarmsIds.isEmpty()) {
if (subscription.isCacheFull()) {
fetchActiveAlarms(subscription);
if (alarmsIds.isEmpty()) {
sendUpdate(subscription.getSessionId(), subscription.createUpdate());
}
} else {
sendUpdate(subscription.getSessionId(), subscription.createUpdate());
}
}
}
} else if (subscription.matches(alarm)) {
if (alarmsIds.size() < alarmsPerAlarmStatusSubscriptionCacheSize) {
alarmsIds.add(alarm.getId().getId());
if (alarmsIds.size() == 1) {
sendUpdate(subscription.getSessionId(), subscription.createUpdate());
}
} else {
subscription.setCacheFull(true);
}
}
} catch (Exception e) {
log.error("[{}, subId: {}] Failed to handle update for alarm status subscription: {}", subscription.getSessionId(), subscription.getSubscriptionId(), subscriptionUpdate, e);
}
}
private boolean validate(TbAbstractSubCtx<?> finalCtx) {
if (finalCtx.isStopped()) {
log.warn("[{}][{}][{}] Received validation task for already stopped context.", finalCtx.getTenantId(), finalCtx.getSessionId(), finalCtx.getCmdId());

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

67
application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmStatusSubscription.java

@ -0,0 +1,67 @@
/**
* 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.AlarmInfo;
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<AlarmSubscriptionUpdate> {
@Getter
private final Set<UUID> alarmIds = new HashSet<>();
@Getter
@Setter
private boolean cacheFull;
@Getter
private final List<String> typeList;
@Getter
private final List<AlarmSeverity> severityList;
@Builder
public TbAlarmStatusSubscription(String serviceId, String sessionId, int subscriptionId, TenantId tenantId, EntityId entityId,
BiConsumer<TbSubscription<AlarmSubscriptionUpdate>, AlarmSubscriptionUpdate> updateProcessor,
List<String> typeList, List<AlarmSeverity> severityList) {
super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.ALARM_STATUS, updateProcessor);
this.typeList = typeList;
this.severityList = severityList;
}
public AlarmStatusUpdate createUpdate() {
return AlarmStatusUpdate.builder()
.cmdId(getSubscriptionId())
.active(alarmIds.size() > 0)
.build();
}
public boolean matches(AlarmInfo alarm) {
return !alarm.isCleared() && (this.typeList == null || this.typeList.contains(alarm.getType())) &&
(this.severityList == null || this.severityList.contains(alarm.getSeverity()));
}
}

3
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);

2
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;

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

8
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.

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

2
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),

42
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<String> typeList;
private List<AlarmSeverity> severityList;
@Override
public WsCmdType getType() {
return WsCmdType.ALARM_STATUS;
}
}

54
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 active;
public AlarmStatusUpdate(int cmdId, boolean active) {
super(cmdId, SubscriptionErrorCode.NO_ERROR.getCode(), null);
this.active = active;
}
public AlarmStatusUpdate(int cmdId, int errorCode, String errorMsg) {
super(cmdId, errorCode, errorMsg);
}
@Builder
public AlarmStatusUpdate(@JsonProperty("cmdId") int cmdId,
@JsonProperty("present") boolean active,
@JsonProperty("errorCode") int errorCode,
@JsonProperty("errorMsg") String errorMsg) {
super(cmdId, errorCode, errorMsg);
this.active = active;
}
@Override
public CmdUpdateType getCmdUpdateType() {
return CmdUpdateType.ALARM_STATUS;
}
}

1
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

2
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.

126
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,124 @@ 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.isActive());
//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.isActive());
//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.isActive());
// 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.isActive());
//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.isActive());
//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.isActive());
}
@Test
public void testAlarmStatusWsCmdWithMaxAlarmsCacheSize() throws Exception {
loginTenantAdmin();
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.assertFalse(update.isActive());
getWsClient().registerWaitForUpdate();
//create 5+1 alarms
List<Alarm> 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);
}
AlarmStatusUpdate updateAfterAlarmsAdded = JacksonUtil.fromString(getWsClient().waitForReply(), AlarmStatusUpdate.class);
Assert.assertEquals(1, updateAfterAlarmsAdded.getCmdId());
Assert.assertTrue(updateAfterAlarmsAdded.isActive());
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.isActive());
}
@Test
public void testEntityDataLatestWidgetFlow() throws Exception {
List<EntityKey> keys = List.of(new EntityKey(EntityKeyType.TIME_SERIES, "temperature"));

3
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<EntitySubtype> findAlarmTypesByTenantId(TenantId tenantId, PageLink pageLink);
List<UUID> findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, int limit);
}

33
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<String> typeList;
private List<AlarmSeverity> severityList;
}

4
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<Alarm> {
PageData<EntitySubtype> findTenantAlarmTypes(UUID tenantId, PageLink pageLink);
boolean removeAlarmTypesIfNoAlarmsPresent(UUID tenantId, Set<String> types);
List<UUID> findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, int limit);
}

7
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<TenantId, Page
return alarmDao.findTenantAlarmTypes(tenantId.getId(), pageLink);
}
@Override
public List<UUID> findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter originatorAlarmFilter, int limit) {
log.trace("Executing findActiveOriginatorAlarms, tenantId [{}], originatorAlarmFilter [{}]", tenantId, originatorAlarmFilter);
return alarmDao.findActiveOriginatorAlarms(tenantId, originatorAlarmFilter, limit);
}
private Alarm merge(Alarm existing, Alarm alarm) {
if (alarm.getStartTs() > existing.getEndTs()) {
existing.setEndTs(alarm.getStartTs());

10
dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java

@ -404,4 +404,14 @@ public interface AlarmRepository extends JpaRepository<AlarmEntity, UUID> {
@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<String> types);
@Query(value = "SELECT a.id FROM alarm a " +
"WHERE a.originator_id = :originatorId " +
"AND (COALESCE(:alarmTypes) IS NULL OR a.type IN (:alarmTypes)) " +
"AND (COALESCE(:alarmSeverities) IS NULL OR a.severity IN (:alarmSeverities)) " +
"AND (a.cleared = false) ORDER BY id LIMIT :limit", nativeQuery = true)
List<UUID> findActiveOriginatorAlarms(@Param("originatorId") UUID originatorId,
@Param("alarmTypes") List<String> alarmTypes,
@Param("alarmSeverities") List<String> alarmSeverities,
int limit);
}

8
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,13 @@ public class JpaAlarmDao extends JpaAbstractDao<AlarmEntity, Alarm> implements A
return alarmRepository.deleteTypeIfNoAlarmsExist(tenantId, types) > 0;
}
@Override
public List<UUID> findActiveOriginatorAlarms(TenantId tenantId, OriginatorAlarmFilter filter, int limit) {
return alarmRepository.findActiveOriginatorAlarms(filter.getOriginatorId().getId(),
filter.getTypeList(), filter.getSeverityList() != null ? filter.getSeverityList().stream().map(Enum::name).toList() : null,
limit);
}
private static String getPropagationTypes(AlarmPropagationInfo ap) {
String propagateRelationTypes;
if (!CollectionUtils.isEmpty(ap.getPropagateRelationTypes())) {

Loading…
Cancel
Save