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 d342446db3..47a79ca2f4 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 @@ -34,20 +34,26 @@ import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.common.data.query.EntityDataPageLink; import org.thingsboard.server.common.data.query.EntityDataQuery; +import org.thingsboard.server.common.data.query.EntityDataSortOrder; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.TsValue; +import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.entity.EntityService; -import org.thingsboard.server.dao.entityview.EntityViewService; +import org.thingsboard.server.dao.model.ModelConstants; 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.telemetry.DefaultTelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; @@ -63,7 +69,6 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; -import java.util.Comparator; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.LinkedHashSet; @@ -88,7 +93,7 @@ import java.util.stream.Collectors; public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubscriptionService { private static final int DEFAULT_LIMIT = 100; - private final Map> subscriptionsBySessionId = new ConcurrentHashMap<>(); + private final Map> subscriptionsBySessionId = new ConcurrentHashMap<>(); @Autowired private TelemetryWebSocketService wsService; @@ -96,6 +101,9 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Autowired private EntityService entityService; + @Autowired + private AlarmService alarmService; + @Autowired @Lazy private TbLocalSubscriptionService localSubscriptionService; @@ -114,10 +122,12 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Value("${database.ts.type}") private String databaseTsType; - @Value("${server.ws.dynamic_page_link_refresh_interval:6}") + @Value("${server.ws.dynamic_page_link.refresh_interval:6}") private long dynamicPageLinkRefreshInterval; - @Value("${server.ws.dynamic_page_link_refresh_pool_size:1}") + @Value("${server.ws.dynamic_page_link.refresh_pool_size:1}") private int dynamicPageLinkRefreshPoolSize; + @Value("${server.ws.max_entities_per_alarm_subscription:1000}") + private int maxEntitiesPerAlarmSubscription; private ExecutorService wsCallBackExecutor; private boolean tsInSqlDB; @@ -195,6 +205,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc ctx.setData(data); ctx.cancelRefreshTask(); if (ctx.getQuery().getPageLink().isDynamic()) { + //TODO: validate number of dynamic page links against rate limits. Ignore dynamic flag if limit is reached. TbEntityDataSubCtx finalCtx = ctx; ScheduledFuture task = scheduler.scheduleWithFixedDelay( () -> refreshDynamicQuery(tenantId, customerId, finalCtx), @@ -230,6 +241,39 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }, wsCallBackExecutor); } + @Override + public void handleCmd(TelemetryWebSocketSessionRef session, AlarmDataCmd cmd) { + TbAlarmDataSubCtx ctx = getSubCtx(session.getSessionId(), cmd.getCmdId()); + if (ctx == null) { + log.debug("[{}][{}] Creating new alarm subscription using: {}", session.getSessionId(), cmd.getCmdId(), cmd); + ctx = createSubCtx(session, cmd); + } + AlarmDataQuery adq = cmd.getQuery(); + EntityDataSortOrder sortOrder = adq.getPageLink().getSortOrder(); + EntityDataSortOrder entitiesSortOrder; + if (sortOrder == null || sortOrder.getKey().getType().equals(EntityKeyType.ALARM_FIELD)) { + entitiesSortOrder = new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, ModelConstants.CREATED_TIME_PROPERTY)); + } else { + entitiesSortOrder = sortOrder; + } + EntityDataPageLink edpl = new EntityDataPageLink(0, maxEntitiesPerAlarmSubscription, null, entitiesSortOrder); + EntityDataQuery edq = new EntityDataQuery(adq.getEntityFilter(), edpl, adq.getEntityFields(), adq.getLatestValues(), adq.getKeyFilters()); + PageData entitiesData = entityService.findEntityDataByQuery(ctx.getTenantId(), ctx.getCustomerId(), edq); + List entities = entitiesData.getData(); + ctx.setEntitiesData(entitiesData); + if (entities.isEmpty()) { + AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(Collections.emptyList(), 1, 0, false), null); + wsService.sendWsMsg(ctx.getSessionId(), update); + } else { + PageData alarms = alarmService.findAlarmDataByQueryForEntities(ctx.getTenantId(), ctx.getCustomerId(), + ctx.getQuery().getPageLink(), ctx.getOrderedEntityIds()); + ctx.setAlarmsData(alarms); + AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), alarms, null); + wsService.sendWsMsg(ctx.getSessionId(), update); + //TODO: Create WS subscription for alarms for this entities. If this is first page(?!) and new alarm matches the filter - invalidate alarms. + } + } + private void refreshDynamicQuery(TenantId tenantId, CustomerId customerId, TbEntityDataSubCtx finalCtx) { try { long start = System.currentTimeMillis(); @@ -244,7 +288,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } - @Scheduled(fixedDelayString = "${server.ws.dynamic_page_link_stats:10000}") + @Scheduled(fixedDelayString = "${server.ws.dynamic_page_link.stats:10000}") public void printStats() { int regularQueryInvocationCntValue = regularQueryInvocationCnt.getAndSet(0); long regularQueryInvocationTimeValue = regularQueryTimeSpent.getAndSet(0); @@ -263,17 +307,25 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { - Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); + Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, sessionRef, cmd.getCmdId()); ctx.setQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; } - private TbEntityDataSubCtx getSubCtx(String sessionId, int cmdId) { - Map sessionSubs = subscriptionsBySessionId.get(sessionId); + private TbAlarmDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) { + Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); + TbAlarmDataSubCtx ctx = new TbAlarmDataSubCtx(serviceId, wsService, sessionRef, cmd.getCmdId()); + ctx.setQuery(cmd.getQuery()); + sessionSubs.put(cmd.getCmdId(), ctx); + return ctx; + } + + private T getSubCtx(String sessionId, int cmdId) { + Map sessionSubs = subscriptionsBySessionId.get(sessionId); if (sessionSubs != null) { - return sessionSubs.get(cmdId); + return (T) sessionSubs.get(cmdId); } else { return null; } @@ -420,14 +472,6 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc return data.stream().collect(Collectors.toMap(TsKvEntry::getKey, value -> new TsValue(value.getTs(), value.getValueAsString()))); } - private Map> toTsValues(List data) { - Map> results = new HashMap<>(); - for (TsKvEntry tsKvEntry : data) { - results.computeIfAbsent(tsKvEntry.getKey(), k -> new ArrayList<>()).add(new TsValue(tsKvEntry.getTs(), tsKvEntry.getValueAsString())); - } - return results; - } - @Override public void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd cmd) { cleanupAndCancel(getSubCtx(sessionId, cmd.getCmdId())); @@ -442,9 +486,9 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Override public void cancelAllSessionSubscriptions(String sessionId) { - Map sessionSubs = subscriptionsBySessionId.remove(sessionId); + Map sessionSubs = subscriptionsBySessionId.remove(sessionId); if (sessionSubs != null) { - sessionSubs.values().forEach(this::cleanupAndCancel); + sessionSubs.values().stream().filter(sub -> sub instanceof TbEntityDataSubCtx).map(sub -> (TbEntityDataSubCtx) sub).forEach(this::cleanupAndCancel); } } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java new file mode 100644 index 0000000000..b323ab7240 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java @@ -0,0 +1,59 @@ +/** + * Copyright © 2016-2020 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.Data; +import lombok.Getter; +import lombok.Setter; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.query.AbstractDataQuery; +import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; +import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; + +@Slf4j +@Data +public abstract class TbAbstractDataSubCtx { + + protected final String serviceId; + protected final TelemetryWebSocketService wsService; + protected final TelemetryWebSocketSessionRef sessionRef; + protected final int cmdId; + @Getter + @Setter + protected T query; + + public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { + this.serviceId = serviceId; + this.wsService = wsService; + this.sessionRef = sessionRef; + this.cmdId = cmdId; + } + + public String getSessionId() { + return sessionRef.getSessionId(); + } + + public TenantId getTenantId() { + return sessionRef.getSecurityCtx().getTenantId(); + } + + public CustomerId getCustomerId() { + return sessionRef.getSecurityCtx().getCustomerId(); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java new file mode 100644 index 0000000000..ab55b5e892 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java @@ -0,0 +1,63 @@ +/** + * Copyright © 2016-2020 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.Getter; +import lombok.Setter; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataQuery; +import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; +import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; + +import java.util.Collection; +import java.util.LinkedHashMap; +import java.util.List; + +@Slf4j +public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { + + @Getter + @Setter + private final LinkedHashMap entitiesMap; + @Getter + @Setter + private boolean tooManyEntities; + + public TbAlarmDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { + super(serviceId, wsService, sessionRef, cmdId); + this.entitiesMap = new LinkedHashMap<>(); + } + + public void setEntitiesData(PageData entitiesData) { + entitiesMap.clear(); + tooManyEntities = entitiesData.hasNext(); + for (EntityData entityData : entitiesData.getData()) { + entitiesMap.put(entityData.getEntityId(), entityData); + } + } + + public Collection getOrderedEntityIds() { + return entitiesMap.keySet(); + } + + public void setAlarmsData(PageData alarms) { +// TODO: implement + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index e3145b4398..fc81362c54 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -17,6 +17,9 @@ package org.thingsboard.server.service.subscription; import lombok.AllArgsConstructor; import lombok.Data; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; @@ -51,17 +54,13 @@ import java.util.function.Function; import java.util.stream.Collectors; @Slf4j -@Data -public class TbEntityDataSubCtx { +public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { - public static final int MAX_SUBS_PER_CMD = 1024 * 8; - private final String serviceId; - private final TelemetryWebSocketService wsService; - private final TelemetryWebSocketSessionRef sessionRef; - private final int cmdId; - private EntityDataQuery query; + @Getter @Setter private TimeSeriesCmd tsCmd; + @Getter private PageData data; + @Getter @Setter private boolean initialDataSent; private Map subToEntityIdMap; private volatile ScheduledFuture refreshTask; @@ -69,22 +68,7 @@ public class TbEntityDataSubCtx { private LatestValueCmd latestValueCmd; public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { - this.serviceId = serviceId; - this.wsService = wsService; - this.sessionRef = sessionRef; - this.cmdId = cmdId; - } - - public String getSessionId() { - return sessionRef.getSessionId(); - } - - public TenantId getTenantId() { - return sessionRef.getSecurityCtx().getTenantId(); - } - - public CustomerId getCustomerId() { - return sessionRef.getSecurityCtx().getCustomerId(); + super(serviceId, wsService, sessionRef, cmdId); } public void setData(PageData data) { 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 af561c8b5b..b805cf576d 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 @@ -16,6 +16,7 @@ package org.thingsboard.server.service.subscription; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; @@ -23,6 +24,8 @@ public interface TbEntityDataSubscriptionService { void handleCmd(TelemetryWebSocketSessionRef sessionId, EntityDataCmd cmd); + void handleCmd(TelemetryWebSocketSessionRef sessionId, AlarmDataCmd cmd); + void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd subscriptionId); void cancelAllSessionSubscriptions(String sessionId); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java index 07bffedcb7..cdcf1c7491 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java @@ -63,6 +63,9 @@ import org.thingsboard.server.service.telemetry.cmd.v1.SubscriptionCmd; import org.thingsboard.server.service.telemetry.cmd.v1.TelemetryPluginCmd; import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.telemetry.cmd.v1.TimeseriesSubscriptionCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.DataCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; @@ -210,6 +213,9 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi if (cmdsWrapper.getEntityDataCmds() != null) { cmdsWrapper.getEntityDataCmds().forEach(cmd -> handleWsEntityDataCmd(sessionRef, cmd)); } + if (cmdsWrapper.getAlarmDataCmds() != null) { + cmdsWrapper.getAlarmDataCmds().forEach(cmd -> handleWsAlarmDataCmd(sessionRef, cmd)); + } if (cmdsWrapper.getEntityDataUnsubscribeCmds() != null) { cmdsWrapper.getEntityDataUnsubscribeCmds().forEach(cmd -> handleWsEntityDataUnsubscribeCmd(sessionRef, cmd)); } @@ -231,6 +237,16 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } } + private void handleWsAlarmDataCmd(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) { + String sessionId = sessionRef.getSessionId(); + log.debug("[{}] Processing: {}", sessionId, cmd); + + if (validateSessionMetadata(sessionRef, cmd.getCmdId(), sessionId) + && validateSubscriptionCmd(sessionRef, cmd)) { + entityDataSubService.handleCmd(sessionRef, cmd); + } + } + private void handleWsEntityDataUnsubscribeCmd(TelemetryWebSocketSessionRef sessionRef, EntityDataUnsubscribeCmd cmd) { String sessionId = sessionRef.getSessionId(); log.debug("[{}] Processing: {}", sessionId, cmd); @@ -246,7 +262,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } @Override - public void sendWsMsg(String sessionId, EntityDataUpdate update) { + public void sendWsMsg(String sessionId, DataUpdate update) { sendWsMsg(sessionId, update.getCmdId(), update); } @@ -661,6 +677,21 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi return true; } + private boolean validateSubscriptionCmd(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) { + if (cmd.getCmdId() < 0) { + SubscriptionUpdate update = new SubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, + "Cmd id is negative value!"); + sendWsMsg(sessionRef, update); + return false; + } else if (cmd.getQuery() == null) { + SubscriptionUpdate update = new SubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, + "Query is empty!"); + sendWsMsg(sessionRef, update); + return false; + } + return true; + } + private boolean validateSubscriptionCmd(TelemetryWebSocketSessionRef sessionRef, SubscriptionCmd cmd) { if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty()) { SubscriptionUpdate update = new SubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java index d04ff71546..896e224898 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.telemetry; +import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; @@ -29,6 +30,6 @@ public interface TelemetryWebSocketService { void sendWsMsg(String sessionId, SubscriptionUpdate update); - void sendWsMsg(String sessionId, EntityDataUpdate update); + void sendWsMsg(String sessionId, DataUpdate update); } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/TelemetryPluginCmdsWrapper.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/TelemetryPluginCmdsWrapper.java index 2eca0a97cf..06842ea7f4 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/TelemetryPluginCmdsWrapper.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/TelemetryPluginCmdsWrapper.java @@ -19,6 +19,8 @@ import lombok.Data; import org.thingsboard.server.service.telemetry.cmd.v1.AttributesSubscriptionCmd; import org.thingsboard.server.service.telemetry.cmd.v1.GetHistoryCmd; import org.thingsboard.server.service.telemetry.cmd.v1.TimeseriesSubscriptionCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; @@ -40,4 +42,8 @@ public class TelemetryPluginCmdsWrapper { private List entityDataUnsubscribeCmds; + private List alarmDataCmds; + + private List alarmDataUnsubscribeCmds; + } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataCmd.java new file mode 100644 index 0000000000..3e9c4bd6f6 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataCmd.java @@ -0,0 +1,33 @@ +/** + * Copyright © 2016-2020 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.telemetry.cmd.v2; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Getter; +import org.thingsboard.server.common.data.query.AlarmDataQuery; + +public class AlarmDataCmd extends DataCmd { + + @Getter + private final AlarmDataQuery query; + + @JsonCreator + public AlarmDataCmd(@JsonProperty("cmdId") int cmdId, @JsonProperty("query") AlarmDataQuery query) { + super(cmdId); + this.query = query; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java new file mode 100644 index 0000000000..6782c2f595 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java @@ -0,0 +1,25 @@ +/** + * Copyright © 2016-2020 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.telemetry.cmd.v2; + +import lombok.Data; + +@Data +public class AlarmDataUnsubscribeCmd { + + private final int cmdId; + +} diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java new file mode 100644 index 0000000000..eac0082db9 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java @@ -0,0 +1,46 @@ +/** + * Copyright © 2016-2020 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.telemetry.cmd.v2; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; + +import java.util.List; + +public class AlarmDataUpdate extends DataUpdate { + + public AlarmDataUpdate(int cmdId, PageData data, List update) { + super(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); + } + + public AlarmDataUpdate(int cmdId, int errorCode, String errorMsg) { + super(cmdId, null, null, errorCode, errorMsg); + } + + @JsonCreator + public AlarmDataUpdate(@JsonProperty("cmdId") int cmdId, + @JsonProperty("data") PageData data, + @JsonProperty("update") List update, + @JsonProperty("errorCode") int errorCode, + @JsonProperty("errorMsg") String errorMsg) { + super(cmdId, data, update, errorCode, errorMsg); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataCmd.java new file mode 100644 index 0000000000..d06e9f107d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataCmd.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2020 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.telemetry.cmd.v2; + +import lombok.Data; +import lombok.Getter; +import lombok.NoArgsConstructor; + +@Data +public class DataCmd { + + @Getter + private final int cmdId; + + public DataCmd(int cmdId) { + this.cmdId = cmdId; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java new file mode 100644 index 0000000000..d35d2bcd53 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2020 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.telemetry.cmd.v2; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; + +import java.util.List; + +@Data +@AllArgsConstructor +public abstract class DataUpdate { + + private final int cmdId; + private final PageData data; + private final List update; + private final int errorCode; + private final String errorMsg; + + public DataUpdate(int cmdId, PageData data, List update) { + this(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); + } + + public DataUpdate(int cmdId, int errorCode, String errorMsg) { + this(cmdId, null, null, errorCode, errorMsg); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataCmd.java index a75451c2c2..7b1f6fb60c 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataCmd.java @@ -15,16 +15,32 @@ */ package org.thingsboard.server.service.telemetry.cmd.v2; -import lombok.Data; +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Getter; import org.thingsboard.server.common.data.query.EntityDataQuery; -@Data -public class EntityDataCmd { +public class EntityDataCmd extends DataCmd { - private final int cmdId; + @Getter private final EntityDataQuery query; + @Getter private final EntityHistoryCmd historyCmd; + @Getter private final LatestValueCmd latestCmd; + @Getter private final TimeSeriesCmd tsCmd; + @JsonCreator + public EntityDataCmd(@JsonProperty("cmdId") int cmdId, + @JsonProperty("query") EntityDataQuery query, + @JsonProperty("historyCmd") EntityHistoryCmd historyCmd, + @JsonProperty("latestCmd") LatestValueCmd latestCmd, + @JsonProperty("tsCmd") TimeSeriesCmd tsCmd) { + super(cmdId); + this.query = query; + this.historyCmd = historyCmd; + this.latestCmd = latestCmd; + this.tsCmd = tsCmd; + } } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java index 3bd51405f8..6f878c42fe 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java @@ -15,30 +15,31 @@ */ package org.thingsboard.server.service.telemetry.cmd.v2; -import lombok.AllArgsConstructor; -import lombok.Data; +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; import java.util.List; -@Data -@AllArgsConstructor -public class EntityDataUpdate { - - private final int cmdId; - private final PageData data; - private final List update; - private final int errorCode; - private final String errorMsg; +public class EntityDataUpdate extends DataUpdate { public EntityDataUpdate(int cmdId, PageData data, List update) { - this(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); + super(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); } public EntityDataUpdate(int cmdId, int errorCode, String errorMsg) { - this(cmdId, null, null, errorCode, errorMsg); + super(cmdId, null, null, errorCode, errorMsg); + } + + @JsonCreator + public EntityDataUpdate(@JsonProperty("cmdId") int cmdId, + @JsonProperty("data") PageData data, + @JsonProperty("update") List update, + @JsonProperty("errorCode") int errorCode, + @JsonProperty("errorMsg") String errorMsg) { + super(cmdId, data, update, errorCode, errorMsg); } } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index b9070a8931..bab6ce6d87 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -46,8 +46,11 @@ server: max_subscriptions_per_regular_user: "${TB_SERVER_WS_TENANT_RATE_LIMITS_MAX_SUBSCRIPTIONS_PER_REGULAR_USER:0}" max_subscriptions_per_public_user: "${TB_SERVER_WS_TENANT_RATE_LIMITS_MAX_SUBSCRIPTIONS_PER_PUBLIC_USER:0}" max_updates_per_session: "${TB_SERVER_WS_TENANT_RATE_LIMITS_MAX_UPDATES_PER_SESSION:300:1,3000:60}" - dynamic_page_link_refresh_interval: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_INTERVAL_SEC:6}" - dynamic_page_link_refresh_pool_size: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_POOL_SIZE:1}" + dynamic_page_link: + refresh_interval: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_INTERVAL_SEC:60}" + refresh_pool_size: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_POOL_SIZE:1}" + max_per_user: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_MAX_PER_USER:10}" + max_entities_per_alarm_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_ALARM_SUBSCRIPTION:1000}" rest: limits: tenant: diff --git a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java index 15da972cf5..0a5dae47b7 100644 --- a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java @@ -26,9 +26,9 @@ import java.util.Arrays; @RunWith(ClasspathSuite.class) @ClasspathSuite.ClassnameFilters({ -// "org.thingsboard.server.controller.sql.WebsocketApiSqlTest", + "org.thingsboard.server.controller.sql.WebsocketApiSqlTest", // "org.thingsboard.server.controller.sql.EntityQueryControllerSqlTest", - "org.thingsboard.server.controller.sql.*Test", +// "org.thingsboard.server.controller.sql.*Test", }) public class ControllerSqlTestSuite { 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 298800a2f0..70dd73c390 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 @@ -24,9 +24,15 @@ import org.thingsboard.server.common.data.alarm.AlarmQuery; 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.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; + +import java.util.Collection; /** * Created by ashvayka on 11.05.17. @@ -52,4 +58,6 @@ public interface AlarmService { ListenableFuture findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type); + PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, + AlarmDataPageLink pageLink, Collection orderedEntityIds); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java new file mode 100644 index 0000000000..36a5aa37ad --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java @@ -0,0 +1,56 @@ +/** + * Copyright © 2016-2020 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 com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.Getter; +import lombok.ToString; + +import java.util.List; + +@ToString +public abstract class AbstractDataQuery extends EntityCountQuery { + + @Getter + protected T pageLink; + @Getter + protected List entityFields; + @Getter + protected List latestValues; + @Getter + protected List keyFilters; + + public AbstractDataQuery() { + super(); + } + + public AbstractDataQuery(EntityFilter entityFilter) { + super(entityFilter); + } + + public AbstractDataQuery(EntityFilter entityFilter, + T pageLink, + List entityFields, + List latestValues, + List keyFilters) { + super(entityFilter); + this.pageLink = pageLink; + this.entityFields = entityFields; + this.latestValues = latestValues; + this.keyFilters = keyFilters; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmData.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmData.java new file mode 100644 index 0000000000..47520100af --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmData.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2020 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.Data; +import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.alarm.AlarmInfo; +import org.thingsboard.server.common.data.id.EntityId; + +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; + +public class AlarmData extends AlarmInfo { + + private final UUID entityId; + private final Map> latest; + + public AlarmData(Alarm alarm, String originatorName, UUID entityId) { + super(alarm, originatorName); + this.entityId = entityId; + this.latest = new HashMap<>(); + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataPageLink.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataPageLink.java new file mode 100644 index 0000000000..589cb2df2f --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataPageLink.java @@ -0,0 +1,67 @@ +/** + * Copyright © 2016-2020 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 com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.Getter; +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 java.util.List; + +@Data +@AllArgsConstructor +public class AlarmDataPageLink extends EntityDataPageLink { + + private long startTs; + private long endTs; + //TODO: handle this; + private long timeWindow; + private List typeList; + private List statusList; + private List severityList; + private boolean searchPropagatedAlarms; + + public AlarmDataPageLink() { + super(); + } + + public AlarmDataPageLink(int pageSize, int page, String textSearch, EntityDataSortOrder sortOrder, boolean dynamic, + boolean searchPropagatedAlarms, + long startTs, long endTs, long timeWindow, + List typeList, List statusList, List severityList) { + super(pageSize, page, textSearch, sortOrder, dynamic); + this.searchPropagatedAlarms = searchPropagatedAlarms; + this.startTs = startTs; + this.endTs = endTs; + this.timeWindow = timeWindow; + this.typeList = typeList; + this.statusList = statusList; + this.severityList = severityList; + } + + @JsonIgnore + public AlarmDataPageLink nextPageLink() { + return new AlarmDataPageLink(this.getPageSize(), this.getPage() + 1, this.getTextSearch(), this.getSortOrder(), this.isDynamic(), + this.searchPropagatedAlarms, + this.startTs, this.endTs, this.timeWindow, + this.typeList, this.statusList, this.severityList + ); + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java new file mode 100644 index 0000000000..1cf9044085 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2020 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 com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.Getter; +import lombok.ToString; + +import java.util.List; + +@ToString +public class AlarmDataQuery extends AbstractDataQuery { + + public AlarmDataQuery() { + } + + public AlarmDataQuery(EntityFilter entityFilter) { + super(entityFilter); + } + + public AlarmDataQuery(EntityFilter entityFilter, AlarmDataPageLink pageLink, List entityFields, List latestValues, List keyFilters) { + super(entityFilter, pageLink, entityFields, latestValues, keyFilters); + } + + @JsonIgnore + public AlarmDataQuery next() { + return new AlarmDataQuery(getEntityFilter(), getPageLink().nextPageLink(), entityFields, latestValues, keyFilters); + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataPageLink.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataPageLink.java index e13fe65158..bdf459ce1f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataPageLink.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataPageLink.java @@ -38,6 +38,6 @@ public class EntityDataPageLink { @JsonIgnore public EntityDataPageLink nextPageLink() { - return new EntityDataPageLink(this.pageSize, this.page+1, this.textSearch, this.sortOrder); + return new EntityDataPageLink(this.pageSize, this.page + 1, this.textSearch, this.sortOrder); } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java index e0c1d12699..41b1551454 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java @@ -22,39 +22,22 @@ import lombok.ToString; import java.util.List; @ToString -public class EntityDataQuery extends EntityCountQuery { - - @Getter - private EntityDataPageLink pageLink; - @Getter - private List entityFields; - @Getter - private List latestValues; - @Getter - private List keyFilters; +public class EntityDataQuery extends AbstractDataQuery { public EntityDataQuery() { - super(); } public EntityDataQuery(EntityFilter entityFilter) { super(entityFilter); } - public EntityDataQuery(EntityFilter entityFilter, - EntityDataPageLink pageLink, - List entityFields, - List latestValues, - List keyFilters) { - super(entityFilter); - this.pageLink = pageLink; - this.entityFields = entityFields; - this.latestValues = latestValues; - this.keyFilters = keyFilters; + public EntityDataQuery(EntityFilter entityFilter, EntityDataPageLink pageLink, List entityFields, List latestValues, List keyFilters) { + super(entityFilter, pageLink, entityFields, latestValues, keyFilters); } @JsonIgnore public EntityDataQuery next() { - return new EntityDataQuery(getEntityFilter(), pageLink.nextPageLink(), entityFields, latestValues, keyFilters); + return new EntityDataQuery(getEntityFilter(), getPageLink().nextPageLink(), entityFields, latestValues, keyFilters); } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKeyType.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKeyType.java index ad8c36f02c..9087927f8c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKeyType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKeyType.java @@ -21,5 +21,6 @@ public enum EntityKeyType { SHARED_ATTRIBUTE, SERVER_ATTRIBUTE, TIME_SERIES, - ENTITY_FIELD; + ENTITY_FIELD, + ALARM_FIELD; } 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 a7eaa3f6d5..ab01ba341d 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 @@ -19,11 +19,16 @@ import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmQuery; +import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.dao.Dao; +import java.util.Collection; import java.util.UUID; /** @@ -40,4 +45,7 @@ public interface AlarmDao extends Dao { Alarm save(TenantId tenantId, Alarm alarm); PageData findAlarms(TenantId tenantId, AlarmQuery query); + + PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, + AlarmDataPageLink pageLink, Collection orderedEntityIds); } 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 6ca6b986f4..c604eb09e6 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 @@ -35,10 +35,14 @@ 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.id.AlarmId; +import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.EntitySearchDirection; @@ -54,6 +58,7 @@ import javax.annotation.Nullable; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.ArrayList; +import java.util.Collection; import java.util.Comparator; import java.util.List; import java.util.Set; @@ -69,6 +74,8 @@ import static org.thingsboard.server.dao.service.Validator.validateId; @Slf4j public class BaseAlarmService extends AbstractEntityService implements AlarmService { + public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; + public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId "; public static final String ALARM_RELATION_PREFIX = "ALARM_"; @Autowired @@ -123,6 +130,14 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ return alarmDao.findLatestByOriginatorAndType(tenantId, originator, type); } + @Override + public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, + AlarmDataPageLink pageLink, Collection orderedEntityIds) { + validateId(tenantId, INCORRECT_TENANT_ID + tenantId); + validateId(customerId, INCORRECT_CUSTOMER_ID + customerId); + return alarmDao.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds); + } + @Override public Boolean deleteAlarm(TenantId tenantId, AlarmId alarmId) { try { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 0f36382d22..08ae361d50 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -227,6 +227,7 @@ public class ModelConstants { public static final String ALARM_TYPE_PROPERTY = "type"; public static final String ALARM_DETAILS_PROPERTY = "details"; public static final String ALARM_ORIGINATOR_ID_PROPERTY = "originator_id"; + public static final String ALARM_ORIGINATOR_NAME_PROPERTY = "originator_name"; public static final String ALARM_ORIGINATOR_TYPE_PROPERTY = "originator_type"; public static final String ALARM_SEVERITY_PROPERTY = "severity"; public static final String ALARM_STATUS_PROPERTY = "status"; 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 f393cac046..8fd5113ed1 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 @@ -20,10 +20,17 @@ import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.query.Param; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.dao.model.sql.AlarmEntity; import org.thingsboard.server.dao.model.sql.AlarmInfoEntity; import org.thingsboard.server.dao.util.SqlDao; +import java.util.Collection; import java.util.List; import java.util.UUID; @@ -72,4 +79,6 @@ public interface AlarmRepository extends CrudRepository { @Param("endTime") Long endTime, @Param("searchText") String searchText, 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 72f07b2085..5bafe2b0e0 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 @@ -26,17 +26,23 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmQuery; import org.thingsboard.server.common.data.alarm.AlarmSearchStatus; +import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.alarm.AlarmDao; import org.thingsboard.server.dao.alarm.BaseAlarmService; import org.thingsboard.server.dao.model.sql.AlarmEntity; import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.sql.JpaAbstractDao; +import org.thingsboard.server.dao.sql.query.AlarmQueryRepository; import org.thingsboard.server.dao.util.SqlDao; +import java.util.Collection; import java.util.List; import java.util.Objects; import java.util.UUID; @@ -55,6 +61,9 @@ public class JpaAlarmDao extends JpaAbstractDao implements A @Autowired private AlarmRepository alarmRepository; + @Autowired + private AlarmQueryRepository alarmQueryRepository; + @Autowired private RelationDao relationDao; @@ -116,4 +125,9 @@ public class JpaAlarmDao extends JpaAbstractDao implements A ) ); } + + @Override + public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataPageLink pageLink, Collection orderedEntityIds) { + return alarmQueryRepository.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds); + } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java new file mode 100644 index 0000000000..582e4fd728 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java @@ -0,0 +1,100 @@ +/** + * Copyright © 2016-2020 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.dao.sql.query; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; +import org.springframework.util.StringUtils; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.EntityIdFactory; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.EntityDataPageLink; +import org.thingsboard.server.dao.model.ModelConstants; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Collectors; + +@Slf4j +public class AlarmDataAdapter { + + private final static ObjectMapper mapper = new ObjectMapper(); + + public static PageData createAlarmData(EntityDataPageLink pageLink, + List> rows, + int totalElements) { + int totalPages = pageLink.getPageSize() > 0 ? (int) Math.ceil((float) totalElements / pageLink.getPageSize()) : 1; + int startIndex = pageLink.getPageSize() * pageLink.getPage(); + boolean hasNext = pageLink.getPageSize() > 0 && totalElements > startIndex + rows.size(); + List entitiesData = convertListToAlarmData(rows); + return new PageData<>(entitiesData, totalPages, totalElements, hasNext); + } + + private static List convertListToAlarmData(List> result) { + return result.stream().map(AlarmDataAdapter::toEntityData).collect(Collectors.toList()); + } + + private static AlarmData toEntityData(Map row) { + Alarm alarm = new Alarm(); + alarm.setId(new AlarmId((UUID) row.get(ModelConstants.ID_PROPERTY))); + alarm.setCreatedTime((long) row.get(ModelConstants.CREATED_TIME_PROPERTY)); + alarm.setAckTs((long) row.get(ModelConstants.ALARM_ACK_TS_PROPERTY)); + alarm.setClearTs((long) row.get(ModelConstants.ALARM_CLEAR_TS_PROPERTY)); + alarm.setStartTs((long) row.get(ModelConstants.ALARM_START_TS_PROPERTY)); + alarm.setEndTs((long) row.get(ModelConstants.ALARM_END_TS_PROPERTY)); + Object additionalInfo = row.get(ModelConstants.ADDITIONAL_INFO_PROPERTY); + if (additionalInfo != null) { + try { + alarm.setDetails(mapper.readTree(additionalInfo.toString())); + } catch (JsonProcessingException e) { + log.warn("Failed to parse json: {}", row.get(ModelConstants.ADDITIONAL_INFO_PROPERTY), e); + } + } + EntityType originatorType = EntityType.values()[(int) row.get(ModelConstants.ALARM_ORIGINATOR_TYPE_PROPERTY)]; + UUID originatorId = (UUID) row.get(ModelConstants.ALARM_ORIGINATOR_ID_PROPERTY); + alarm.setOriginator(EntityIdFactory.getByTypeAndUuid(originatorType, originatorId)); + alarm.setPropagate((boolean) row.get(ModelConstants.ALARM_PROPAGATE_PROPERTY)); + alarm.setType(row.get(ModelConstants.ALARM_TYPE_PROPERTY).toString()); + alarm.setSeverity(AlarmSeverity.valueOf(row.get(ModelConstants.ALARM_SEVERITY_PROPERTY).toString())); + alarm.setStatus(AlarmStatus.valueOf(row.get(ModelConstants.ALARM_STATUS_PROPERTY).toString())); + alarm.setTenantId(new TenantId((UUID) row.get(ModelConstants.TENANT_ID_PROPERTY))); + if (row.get(ModelConstants.ALARM_PROPAGATE_RELATION_TYPES) != null) { + String propagateRelationTypes = row.get(ModelConstants.ALARM_PROPAGATE_RELATION_TYPES).toString(); + if (!StringUtils.isEmpty(propagateRelationTypes)) { + alarm.setPropagateRelationTypes(Arrays.asList(propagateRelationTypes.split(","))); + } else { + alarm.setPropagateRelationTypes(Collections.emptyList()); + } + } else { + alarm.setPropagateRelationTypes(Collections.emptyList()); + } + UUID entityId = (UUID) row.get(ModelConstants.ENTITY_ID_COLUMN); + Object originatorNameObj = row.get(ModelConstants.ALARM_ORIGINATOR_NAME_PROPERTY); + String originatorName = originatorNameObj != null ? originatorNameObj.toString() : null; + return new AlarmData(alarm, originatorName, entityId); + } + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java new file mode 100644 index 0000000000..25f8d98aac --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2020 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.dao.sql.query; + +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; + +import java.util.Collection; + +public interface AlarmQueryRepository { + + PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, + AlarmDataPageLink pageLink, Collection orderedEntityIds); + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java new file mode 100644 index 0000000000..cd391b01b9 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java @@ -0,0 +1,241 @@ +/** + * Copyright © 2016-2020 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.dao.sql.query; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; +import org.springframework.stereotype.Repository; +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.id.CustomerId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.EntityDataSortOrder; +import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.dao.model.ModelConstants; +import org.thingsboard.server.dao.util.SqlDao; + +import java.util.Collection; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; + +@SqlDao +@Repository +@Slf4j +public class DefaultAlarmQueryRepository implements AlarmQueryRepository { + + private static final Map alarmFieldColumnMap = new HashMap<>(); + + static { + alarmFieldColumnMap.put("createdTime", ModelConstants.CREATED_TIME_PROPERTY); + alarmFieldColumnMap.put("ackTs", ModelConstants.ALARM_ACK_TS_PROPERTY); + alarmFieldColumnMap.put("clearTs", ModelConstants.ALARM_CLEAR_TS_PROPERTY); + alarmFieldColumnMap.put("details", ModelConstants.ADDITIONAL_INFO_PROPERTY); + alarmFieldColumnMap.put("endTs", ModelConstants.ALARM_END_TS_PROPERTY); + alarmFieldColumnMap.put("startTs", ModelConstants.ALARM_START_TS_PROPERTY); + alarmFieldColumnMap.put("status", ModelConstants.ALARM_STATUS_PROPERTY); + alarmFieldColumnMap.put("type", ModelConstants.ALARM_TYPE_PROPERTY); + alarmFieldColumnMap.put("severity", ModelConstants.ALARM_SEVERITY_PROPERTY); + alarmFieldColumnMap.put("originator_id", ModelConstants.ALARM_ORIGINATOR_ID_PROPERTY); + alarmFieldColumnMap.put("originator_type", ModelConstants.ALARM_ORIGINATOR_TYPE_PROPERTY); + } + + public static final String SELECT_ORIGINATOR_NAME = " CASE" + + " WHEN a.originator_type = 0" + + " THEN (select title from tenant where id = a.originator_id)" + + " WHEN a.originator_type = 1 " + + " THEN (select title from customer where id = a.originator_id)" + + " WHEN a.originator_type = 2" + + " THEN (select CONCAT (first_name, ' ', last_name) from tb_user where id = a.originator_id)" + + " WHEN a.originator_type = 3" + + " THEN (select title from dashboard where id = a.originator_id)" + + " WHEN a.originator_type = 4" + + " THEN (select name from asset where id = a.originator_id)" + + " WHEN a.originator_type = 5" + + " THEN (select name from device where id = a.originator_id)" + + " WHEN a.originator_type = 9" + + " THEN (select name from entity_view where id = a.originator_id)" + + " END as originator_name"; + + public static final String FIELDS_SELECTION = "select a.id as id," + + " a.created_time as created_time," + + " a.ack_ts as ack_ts," + + " a.clear_ts as clear_ts," + + " a.additional_info as additional_info," + + " a.end_ts as end_ts," + + " a.originator_id as originator_id," + + " a.originator_type as originator_type," + + " a.propagate as propagate," + + " a.severity as severity," + + " a.start_ts as start_ts," + + " a.status as status, " + + " a.tenant_id as tenant_id, " + + " a.propagate_relation_types as propagate_relation_types, " + + " a.type as type," + SELECT_ORIGINATOR_NAME + ", "; + + public static final String JOIN_RELATIONS = "left join relation r on r.relation_type_group = 'ALARM' and r.relation_type = 'ALARM_ANY' and a.id = r.to_id"; + + @Autowired + protected NamedParameterJdbcTemplate jdbcTemplate; + + + @Override + public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, + AlarmDataPageLink pageLink, Collection orderedEntityIds) { + QueryContext ctx = new QueryContext(); + + StringBuilder selectPart = new StringBuilder(FIELDS_SELECTION); + StringBuilder fromPart = new StringBuilder(" from alarm a "); + StringBuilder wherePart = new StringBuilder(" where "); + StringBuilder sortPart = new StringBuilder(" order by "); + boolean addAnd = false; + if (pageLink.isSearchPropagatedAlarms()) { + selectPart.append(" r.from_id as entity_id "); + fromPart.append(JOIN_RELATIONS); + } else { + selectPart.append(" a.originator_id as entity_id "); + } + EntityDataSortOrder sortOrder = pageLink.getSortOrder(); + if (sortOrder != null && sortOrder.getKey().getType().equals(EntityKeyType.ALARM_FIELD)) { + String sortOrderKey = sortOrder.getKey().getKey(); + sortPart.append("a.").append(alarmFieldColumnMap.getOrDefault(sortOrderKey, sortOrderKey)) + .append(" ").append(sortOrder.getDirection().name()); + ctx.addUuidListParameter("entity_ids", orderedEntityIds.stream().map(EntityId::getId).collect(Collectors.toList())); + if (pageLink.isSearchPropagatedAlarms()) { + fromPart.append(" and r.from_id in (:entity_ids)"); + } else { + wherePart.append(" a.originator_id in (:entity_ids)"); + addAnd = true; + } + } else { + fromPart.append(" left join (select * from (VALUES"); + int entityIdIdx = 0; + int lastEntityIdIdx = orderedEntityIds.size() - 1; + for (EntityId entityId : orderedEntityIds) { + fromPart.append("(uuid('").append(entityId.getId().toString()).append("'), ").append(entityIdIdx).append(")"); + if (entityIdIdx != lastEntityIdIdx) { + fromPart.append(","); + } else { + fromPart.append(")"); + } + entityIdIdx++; + } + fromPart.append(" as e(id, priority)) e "); + if (pageLink.isSearchPropagatedAlarms()) { + fromPart.append("on r.from_id = e.id"); + } else { + fromPart.append("on a.originator_id = e.id"); + } + sortPart.append("e.priority"); + } + + if (pageLink.getStartTs() > 0) { + addAndIfNeeded(wherePart, addAnd); + addAnd = true; + ctx.addLongParameter("startTime", pageLink.getStartTs()); + wherePart.append("a.created_time >= :startTime"); + } + + if (pageLink.getEndTs() > 0) { + addAndIfNeeded(wherePart, addAnd); + addAnd = true; + ctx.addLongParameter("endTime", pageLink.getEndTs()); + wherePart.append("a.created_time <= :endTime"); + } + + if (pageLink.getTypeList() != null && !pageLink.getTypeList().isEmpty()) { + addAndIfNeeded(wherePart, addAnd); + addAnd = true; + ctx.addStringListParameter("alarmTypes", pageLink.getTypeList()); + wherePart.append("a.type in (:alarmTypes)"); + } + + if (pageLink.getSeverityList() != null && !pageLink.getSeverityList().isEmpty()) { + addAndIfNeeded(wherePart, addAnd); + addAnd = true; + ctx.addStringListParameter("alarmSeverities", pageLink.getSeverityList().stream().map(AlarmSeverity::name).collect(Collectors.toList())); + wherePart.append("a.severity in (:alarmSeverities)"); + } + + if (pageLink.getStatusList() != null && !pageLink.getStatusList().isEmpty()) { + Set statusSet = toStatusSet(pageLink.getStatusList()); + if (!statusSet.isEmpty()) { + addAndIfNeeded(wherePart, addAnd); + addAnd = true; + ctx.addStringListParameter("alarmStatuses", statusSet.stream().map(AlarmStatus::name).collect(Collectors.toList())); + wherePart.append(" a.status in (:alarmStatuses)"); + } + } + + String countQuery = fromPart.toString() + wherePart.toString(); + int totalElements = jdbcTemplate.queryForObject(String.format("select count(*) %s", countQuery), ctx, Integer.class); + + String dataQuery = selectPart.toString() + countQuery + sortPart; + + int startIndex = pageLink.getPageSize() * pageLink.getPage(); + if (pageLink.getPageSize() > 0) { + dataQuery = String.format("%s limit %s offset %s", dataQuery, pageLink.getPageSize(), startIndex); + } + List> rows = jdbcTemplate.queryForList(dataQuery, ctx); + return AlarmDataAdapter.createAlarmData(pageLink, rows, totalElements); + } + + private Set toStatusSet(List statusList) { + Set result = new HashSet<>(); + for (AlarmSearchStatus searchStatus : statusList) { + switch (searchStatus) { + case ACK: + result.add(AlarmStatus.ACTIVE_ACK); + result.add(AlarmStatus.CLEARED_ACK); + break; + case UNACK: + result.add(AlarmStatus.ACTIVE_UNACK); + result.add(AlarmStatus.CLEARED_UNACK); + break; + case CLEARED: + result.add(AlarmStatus.CLEARED_ACK); + result.add(AlarmStatus.CLEARED_UNACK); + break; + case ACTIVE: + result.add(AlarmStatus.ACTIVE_ACK); + result.add(AlarmStatus.ACTIVE_UNACK); + break; + default: + break; + } + if (searchStatus == AlarmSearchStatus.ANY || result.size() == AlarmStatus.values().length) { + result.clear(); + return result; + } + } + return result; + } + + private void addAndIfNeeded(StringBuilder wherePart, boolean addAnd) { + if (addAnd) { + wherePart.append(" and "); + } + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java index 5dd5f51032..2416d56d55 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java @@ -218,7 +218,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { @Override public long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query) { EntityType entityType = resolveEntityType(query.getEntityFilter()); - EntityQueryContext ctx = new EntityQueryContext(); + QueryContext ctx = new QueryContext(); ctx.append("select count(e.id) from "); ctx.append(addEntityTableQuery(ctx, query.getEntityFilter(), entityType)); ctx.append(" e where "); @@ -228,7 +228,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { @Override public PageData findEntityDataByQuery(TenantId tenantId, CustomerId customerId, EntityDataQuery query) { - EntityQueryContext ctx = new EntityQueryContext(); + QueryContext ctx = new QueryContext(); EntityType entityType = resolveEntityType(query.getEntityFilter()); EntityDataPageLink pageLink = query.getPageLink(); @@ -308,7 +308,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { return EntityDataAdapter.createEntityData(pageLink, selectionMapping, rows, totalElements); } - private String buildEntityWhere(EntityQueryContext ctx, + private String buildEntityWhere(QueryContext ctx, TenantId tenantId, CustomerId customerId, EntityFilter entityFilter, @@ -327,7 +327,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { return result; } - private String buildPermissionQuery(EntityQueryContext ctx, EntityFilter entityFilter, TenantId tenantId, CustomerId customerId, EntityType entityType) { + private String buildPermissionQuery(QueryContext ctx, EntityFilter entityFilter, TenantId tenantId, CustomerId customerId, EntityType entityType) { switch (entityFilter.getType()) { case RELATIONS_QUERY: case DEVICE_SEARCH_QUERY: @@ -343,7 +343,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { } } - private String defaultPermissionQuery(EntityQueryContext ctx, TenantId tenantId, CustomerId customerId, EntityType entityType) { + private String defaultPermissionQuery(QueryContext ctx, TenantId tenantId, CustomerId customerId, EntityType entityType) { ctx.addUuidParameter("permissions_tenant_id", tenantId.getId()); if (customerId != null && !customerId.isNullUid()) { ctx.addUuidParameter("permissions_customer_id", customerId.getId()); @@ -357,7 +357,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { } } - private String buildEntityFilterQuery(EntityQueryContext ctx, EntityFilter entityFilter) { + private String buildEntityFilterQuery(QueryContext ctx, EntityFilter entityFilter) { switch (entityFilter.getType()) { case SINGLE_ENTITY: return this.singleEntityQuery(ctx, (SingleEntityFilter) entityFilter); @@ -378,7 +378,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { } } - private String addEntityTableQuery(EntityQueryContext ctx, EntityFilter entityFilter, EntityType entityType) { + private String addEntityTableQuery(QueryContext ctx, EntityFilter entityFilter, EntityType entityType) { switch (entityFilter.getType()) { case RELATIONS_QUERY: return relationQuery(ctx, (RelationsQueryFilter) entityFilter); @@ -393,7 +393,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { } } - private String entitySearchQuery(EntityQueryContext ctx, EntitySearchQueryFilter entityFilter, EntityType entityType, List types) { + private String entitySearchQuery(QueryContext ctx, EntitySearchQueryFilter entityFilter, EntityType entityType, List types) { EntityId rootId = entityFilter.getRootEntity(); //TODO: fetch last level only. //TODO: fetch distinct records. @@ -416,7 +416,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { return query; } - private String relationQuery(EntityQueryContext ctx, RelationsQueryFilter entityFilter) { + private String relationQuery(QueryContext ctx, RelationsQueryFilter entityFilter) { EntityId rootId = entityFilter.getRootEntity(); String lvlFilter = getLvlFilter(entityFilter.getMaxLevel()); String selectFields = SELECT_TENANT_ID + ", " + SELECT_CUSTOMER_ID @@ -478,7 +478,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { return from; } - private String buildWhere(EntityQueryContext ctx, List latestFiltersMapping) { + private String buildWhere(QueryContext ctx, List latestFiltersMapping) { String latestFilters = EntityKeyMapping.buildQuery(ctx, latestFiltersMapping); if (!StringUtils.isEmpty(latestFilters)) { return String.format("where %s", latestFilters); @@ -487,7 +487,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { } } - private String buildTextSearchQuery(EntityQueryContext ctx, List selectionMapping, String searchText) { + private String buildTextSearchQuery(QueryContext ctx, List selectionMapping, String searchText) { if (!StringUtils.isEmpty(searchText) && !selectionMapping.isEmpty()) { String lowerSearchText = searchText.toLowerCase() + "%"; List searchPredicates = selectionMapping.stream().map(mapping -> { @@ -502,22 +502,22 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { } } - private String singleEntityQuery(EntityQueryContext ctx, SingleEntityFilter filter) { + private String singleEntityQuery(QueryContext ctx, SingleEntityFilter filter) { ctx.addUuidParameter("entity_filter_single_entity_id", filter.getSingleEntity().getId()); return "e.id=:entity_filter_single_entity_id"; } - private String entityListQuery(EntityQueryContext ctx, EntityListFilter filter) { + private String entityListQuery(QueryContext ctx, EntityListFilter filter) { ctx.addUuidListParameter("entity_filter_entity_ids", filter.getEntityList().stream().map(UUID::fromString).collect(Collectors.toList())); return "e.id in (:entity_filter_entity_ids)"; } - private String entityNameQuery(EntityQueryContext ctx, EntityNameFilter filter) { + private String entityNameQuery(QueryContext ctx, EntityNameFilter filter) { ctx.addStringParameter("entity_filter_name_filter", filter.getEntityNameFilter()); return "lower(e.search_text) like lower(concat(:entity_filter_name_filter, '%%'))"; } - private String typeQuery(EntityQueryContext ctx, EntityFilter filter) { + private String typeQuery(QueryContext ctx, EntityFilter filter) { String type; String name; switch (filter.getType()) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java index 4fcdf0f777..1b91178fd7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java @@ -194,7 +194,7 @@ public class EntityKeyMapping { } } - public Stream toQueries(EntityQueryContext ctx) { + public Stream toQueries(QueryContext ctx) { if (hasFilter()) { String keyAlias = entityKey.getType().equals(EntityKeyType.ENTITY_FIELD) ? "e" : alias; return keyFilters.stream().map(keyFilter -> @@ -204,7 +204,7 @@ public class EntityKeyMapping { } } - public String toLatestJoin(EntityQueryContext ctx, EntityFilter entityFilter, EntityType entityType) { + public String toLatestJoin(QueryContext ctx, EntityFilter entityFilter, EntityType entityType) { String entityTypeStr; if (entityFilter.getType().equals(EntityFilterType.RELATIONS_QUERY)) { entityTypeStr = "entities.entity_type"; @@ -239,12 +239,12 @@ public class EntityKeyMapping { Collectors.joining(", ")); } - public static String buildLatestJoins(EntityQueryContext ctx, EntityFilter entityFilter, EntityType entityType, List latestMappings) { + public static String buildLatestJoins(QueryContext ctx, EntityFilter entityFilter, EntityType entityType, List latestMappings) { return latestMappings.stream().map(mapping -> mapping.toLatestJoin(ctx, entityFilter, entityType)).collect( Collectors.joining(" ")); } - public static String buildQuery(EntityQueryContext ctx, List mappings) { + public static String buildQuery(QueryContext ctx, List mappings) { return mappings.stream().flatMap(mapping -> mapping.toQueries(ctx)).collect( Collectors.joining(" AND ")); } @@ -357,11 +357,11 @@ public class EntityKeyMapping { return String.join(", ", attrValSelection, attrTsSelection); } - private String buildKeyQuery(EntityQueryContext ctx, String alias, KeyFilter keyFilter) { + private String buildKeyQuery(QueryContext ctx, String alias, KeyFilter keyFilter) { return this.buildPredicateQuery(ctx, alias, keyFilter.getKey(), keyFilter.getPredicate()); } - private String buildPredicateQuery(EntityQueryContext ctx, String alias, EntityKey key, KeyFilterPredicate predicate) { + private String buildPredicateQuery(QueryContext ctx, String alias, EntityKey key, KeyFilterPredicate predicate) { if (predicate.getType().equals(FilterPredicateType.COMPLEX)) { return this.buildComplexPredicateQuery(ctx, alias, key, (ComplexFilterPredicate) predicate); } else { @@ -369,14 +369,14 @@ public class EntityKeyMapping { } } - private String buildComplexPredicateQuery(EntityQueryContext ctx, String alias, EntityKey key, ComplexFilterPredicate predicate) { + private String buildComplexPredicateQuery(QueryContext ctx, String alias, EntityKey key, ComplexFilterPredicate predicate) { return predicate.getPredicates().stream() .map(keyFilterPredicate -> this.buildPredicateQuery(ctx, alias, key, keyFilterPredicate)).collect(Collectors.joining( " " + predicate.getOperation().name() + " " )); } - private String buildSimplePredicateQuery(EntityQueryContext ctx, String alias, EntityKey key, KeyFilterPredicate predicate) { + private String buildSimplePredicateQuery(QueryContext ctx, String alias, EntityKey key, KeyFilterPredicate predicate) { if (predicate.getType().equals(FilterPredicateType.NUMERIC)) { if (key.getType().equals(EntityKeyType.ENTITY_FIELD)) { String column = entityFieldColumnMap.get(key.getKey()); @@ -402,7 +402,7 @@ public class EntityKeyMapping { } } - private String buildStringPredicateQuery(EntityQueryContext ctx, String field, StringFilterPredicate stringFilterPredicate) { + private String buildStringPredicateQuery(QueryContext ctx, String field, StringFilterPredicate stringFilterPredicate) { String operationField = field; String paramName = getNextParameterName(field); String value = stringFilterPredicate.getValue(); @@ -439,7 +439,7 @@ public class EntityKeyMapping { return String.format("(%s is not null and %s)", field, stringOperationQuery); } - private String buildNumericPredicateQuery(EntityQueryContext ctx, String field, NumericFilterPredicate numericFilterPredicate) { + private String buildNumericPredicateQuery(QueryContext ctx, String field, NumericFilterPredicate numericFilterPredicate) { String paramName = getNextParameterName(field); ctx.addDoubleParameter(paramName, numericFilterPredicate.getValue()); String numericOperationQuery = ""; @@ -466,7 +466,7 @@ public class EntityKeyMapping { return String.format("(%s is not null and %s)", field, numericOperationQuery); } - private String buildBooleanPredicateQuery(EntityQueryContext ctx, String field, + private String buildBooleanPredicateQuery(QueryContext ctx, String field, BooleanFilterPredicate booleanFilterPredicate) { String paramName = getNextParameterName(field); ctx.addBooleanParameter(paramName, booleanFilterPredicate.isValue()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityQueryContext.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/QueryContext.java similarity index 94% rename from dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityQueryContext.java rename to dao/src/main/java/org/thingsboard/server/dao/sql/query/QueryContext.java index 9aa433c300..4ffaedbbc9 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityQueryContext.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/QueryContext.java @@ -24,13 +24,13 @@ import java.util.List; import java.util.Map; import java.util.UUID; -public class EntityQueryContext implements SqlParameterSource { +public class QueryContext implements SqlParameterSource { private static final PostgresUUIDType UUID_TYPE = new PostgresUUIDType(); private final StringBuilder query; private final Map params; - public EntityQueryContext() { + public QueryContext() { query = new StringBuilder(); params = new HashMap<>(); } @@ -91,6 +91,10 @@ public class EntityQueryContext implements SqlParameterSource { addParameter(name, value, Types.DOUBLE, "DOUBLE"); } + public void addLongParameter(String name, long value) { + addParameter(name, value, Types.BIGINT, "BIGINT"); + } + public void addStringListParameter(String name, List value) { addParameter(name, value, Types.VARCHAR, "VARCHAR"); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java index 5d494546c9..e02befcd84 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java @@ -24,16 +24,26 @@ import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmQuery; +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.id.AssetId; +import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.data.query.AlarmData; +import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; +import org.thingsboard.server.common.data.query.EntityDataSortOrder; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; +import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.concurrent.ExecutionException; @@ -195,6 +205,128 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { Assert.assertEquals(created, alarms.getData().get(0)); } + @Test + public void testFindAlarmUsingAlarmDataQuery() throws ExecutionException, InterruptedException { + AssetId parentId = new AssetId(Uuids.timeBased()); + AssetId childId = new AssetId(Uuids.timeBased()); + + EntityRelation relation = new EntityRelation(parentId, childId, EntityRelation.CONTAINS_TYPE); + + Assert.assertTrue(relationService.saveRelationAsync(tenantId, relation).get()); + + long ts = System.currentTimeMillis(); + Alarm alarm = Alarm.builder().tenantId(tenantId).originator(childId) + .type(TEST_ALARM) + .propagate(false) + .severity(AlarmSeverity.CRITICAL) + .status(AlarmStatus.ACTIVE_UNACK) + .startTs(ts).build(); + + Alarm created = alarmService.createOrUpdateAlarm(alarm); + + AlarmDataPageLink pageLink = new AlarmDataPageLink(); + pageLink.setPage(0); + pageLink.setPageSize(1); + pageLink.setSortOrder(new EntityDataSortOrder(new EntityKey(EntityKeyType.ALARM_FIELD, "createdTime"))); + + pageLink.setStartTs(0L); + pageLink.setEndTs(System.currentTimeMillis()); + pageLink.setSearchPropagatedAlarms(false); + pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); + pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); + + PageData alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(childId)); + + Assert.assertNotNull(alarms.getData()); + Assert.assertEquals(1, alarms.getData().size()); + Assert.assertEquals(created, alarms.getData().get(0)); + + pageLink.setPage(0); + pageLink.setPageSize(1); + pageLink.setSortOrder(new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, "createdTime"))); + + pageLink.setStartTs(0L); + pageLink.setEndTs(System.currentTimeMillis()); + pageLink.setSearchPropagatedAlarms(false); + pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); + pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); + + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(childId)); + Assert.assertNotNull(alarms.getData()); + Assert.assertEquals(1, alarms.getData().size()); + Assert.assertEquals(created, alarms.getData().get(0)); + + // Check child relation + Assert.assertNotNull(alarms.getData()); + Assert.assertEquals(1, alarms.getData().size()); + Assert.assertEquals(created, new Alarm(alarms.getData().get(0))); + + created.setPropagate(true); + created = alarmService.createOrUpdateAlarm(created); + + // Check child relation + pageLink.setPage(0); + pageLink.setPageSize(1); + pageLink.setSortOrder(new EntityDataSortOrder(new EntityKey(EntityKeyType.ALARM_FIELD, "createdTime"))); + + pageLink.setStartTs(0L); + pageLink.setEndTs(System.currentTimeMillis()); + pageLink.setSearchPropagatedAlarms(true); + pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); + pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); + + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(childId)); + + // Check parent relation + pageLink.setPage(0); + pageLink.setPageSize(1); + pageLink.setSortOrder(new EntityDataSortOrder(new EntityKey(EntityKeyType.ALARM_FIELD, "createdTime"))); + + pageLink.setStartTs(0L); + pageLink.setEndTs(System.currentTimeMillis()); + pageLink.setSearchPropagatedAlarms(true); + pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); + pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); + + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(parentId)); + Assert.assertNotNull(alarms.getData()); + Assert.assertEquals(1, alarms.getData().size()); + Assert.assertEquals(created, alarms.getData().get(0)); + + pageLink.setPage(0); + pageLink.setPageSize(1); + pageLink.setSortOrder(new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, "createdTime"))); + + pageLink.setStartTs(0L); + pageLink.setEndTs(System.currentTimeMillis()); + pageLink.setSearchPropagatedAlarms(true); + pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); + pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); + + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(parentId)); + Assert.assertNotNull(alarms.getData()); + Assert.assertEquals(1, alarms.getData().size()); + Assert.assertEquals(created, alarms.getData().get(0)); + + alarmService.ackAlarm(tenantId, created.getId(), System.currentTimeMillis()).get(); + created = alarmService.findAlarmByIdAsync(tenantId, created.getId()).get(); + + pageLink.setPage(0); + pageLink.setPageSize(1); + pageLink.setSortOrder(new EntityDataSortOrder(new EntityKey(EntityKeyType.ALARM_FIELD, "createdTime"))); + + pageLink.setStartTs(0L); + pageLink.setEndTs(System.currentTimeMillis()); + pageLink.setSearchPropagatedAlarms(true); + pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); + pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); + + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(childId)); + Assert.assertNotNull(alarms.getData()); + Assert.assertEquals(1, alarms.getData().size()); + Assert.assertEquals(created, alarms.getData().get(0)); + } + @Test public void testDeleteAlarm() throws ExecutionException, InterruptedException { AssetId parentId = new AssetId(Uuids.timeBased()); diff --git a/ui-ngx/package-lock.json b/ui-ngx/package-lock.json index b9d50ac759..7def972724 100644 --- a/ui-ngx/package-lock.json +++ b/ui-ngx/package-lock.json @@ -8997,10 +8997,10 @@ "integrity": "sha512-4O3GWAYJaauMCILm07weko2rHA8a4kjn7+8Lg4s1d7SxwS/3IpkVD/GljbRrIJ1c1W/XGJ3GbuK7RyYZEJChhw==" }, "ngx-flowchart": { - "version": "git://github.com/thingsboard/ngx-flowchart.git#7a02f4748b5e7821a883c903107af5f20415d026", + "version": "git://github.com/thingsboard/ngx-flowchart.git#a4157b0eef2eb3646ef920447c7b06b39d54f87f", "from": "git://github.com/thingsboard/ngx-flowchart.git#master", "requires": { - "tslib": "^1.13.0" + "tslib": "^1.10.0" }, "dependencies": { "tslib": {