Browse Source

Merge branch 'feature/alarm-data-query' of github.com:thingsboard/thingsboard into feature/entity-data-query

pull/3053/head
Andrii Shvaika 6 years ago
parent
commit
134c390cb2
  1. 86
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 59
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java
  3. 63
      application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java
  4. 32
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java
  5. 3
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java
  6. 33
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java
  7. 3
      application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java
  8. 6
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/TelemetryPluginCmdsWrapper.java
  9. 33
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataCmd.java
  10. 25
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java
  11. 46
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java
  12. 32
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataCmd.java
  13. 43
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java
  14. 24
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataCmd.java
  15. 27
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java
  16. 7
      application/src/main/resources/thingsboard.yml
  17. 4
      application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java
  18. 8
      common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java
  19. 56
      common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java
  20. 37
      common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmData.java
  21. 67
      common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataPageLink.java
  22. 42
      common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java
  23. 2
      common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataPageLink.java
  24. 27
      common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java
  25. 3
      common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKeyType.java
  26. 8
      dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java
  27. 15
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
  28. 1
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  29. 9
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java
  30. 14
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java
  31. 100
      dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmDataAdapter.java
  32. 32
      dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java
  33. 241
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java
  34. 30
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java
  35. 22
      dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java
  36. 8
      dao/src/main/java/org/thingsboard/server/dao/sql/query/QueryContext.java
  37. 132
      dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java
  38. 4
      ui-ngx/package-lock.json

86
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.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData; 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.EntityData;
import org.thingsboard.server.common.data.query.EntityDataPageLink;
import org.thingsboard.server.common.data.query.EntityDataQuery; 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.EntityKey;
import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.query.TsValue; 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.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.dao.timeseries.TimeseriesService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.executors.DbCallbackExecutorService; 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.TelemetryWebSocketService;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; 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.EntityDataCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate;
@ -63,7 +69,6 @@ import java.util.ArrayList;
import java.util.Arrays; import java.util.Arrays;
import java.util.Collection; import java.util.Collection;
import java.util.Collections; import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap; import java.util.HashMap;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
import java.util.LinkedHashSet; import java.util.LinkedHashSet;
@ -88,7 +93,7 @@ import java.util.stream.Collectors;
public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubscriptionService { public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubscriptionService {
private static final int DEFAULT_LIMIT = 100; private static final int DEFAULT_LIMIT = 100;
private final Map<String, Map<Integer, TbEntityDataSubCtx>> subscriptionsBySessionId = new ConcurrentHashMap<>(); private final Map<String, Map<Integer, TbAbstractDataSubCtx>> subscriptionsBySessionId = new ConcurrentHashMap<>();
@Autowired @Autowired
private TelemetryWebSocketService wsService; private TelemetryWebSocketService wsService;
@ -96,6 +101,9 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
@Autowired @Autowired
private EntityService entityService; private EntityService entityService;
@Autowired
private AlarmService alarmService;
@Autowired @Autowired
@Lazy @Lazy
private TbLocalSubscriptionService localSubscriptionService; private TbLocalSubscriptionService localSubscriptionService;
@ -114,10 +122,12 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
@Value("${database.ts.type}") @Value("${database.ts.type}")
private String databaseTsType; private String databaseTsType;
@Value("${server.ws.dynamic_page_link_refresh_interval:6}") @Value("${server.ws.dynamic_page_link.refresh_interval:6}")
private long dynamicPageLinkRefreshInterval; 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; private int dynamicPageLinkRefreshPoolSize;
@Value("${server.ws.max_entities_per_alarm_subscription:1000}")
private int maxEntitiesPerAlarmSubscription;
private ExecutorService wsCallBackExecutor; private ExecutorService wsCallBackExecutor;
private boolean tsInSqlDB; private boolean tsInSqlDB;
@ -195,6 +205,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
ctx.setData(data); ctx.setData(data);
ctx.cancelRefreshTask(); ctx.cancelRefreshTask();
if (ctx.getQuery().getPageLink().isDynamic()) { 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; TbEntityDataSubCtx finalCtx = ctx;
ScheduledFuture<?> task = scheduler.scheduleWithFixedDelay( ScheduledFuture<?> task = scheduler.scheduleWithFixedDelay(
() -> refreshDynamicQuery(tenantId, customerId, finalCtx), () -> refreshDynamicQuery(tenantId, customerId, finalCtx),
@ -230,6 +241,39 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
}, wsCallBackExecutor); }, 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<EntityData> entitiesData = entityService.findEntityDataByQuery(ctx.getTenantId(), ctx.getCustomerId(), edq);
List<EntityData> 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<AlarmData> 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) { private void refreshDynamicQuery(TenantId tenantId, CustomerId customerId, TbEntityDataSubCtx finalCtx) {
try { try {
long start = System.currentTimeMillis(); 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() { public void printStats() {
int regularQueryInvocationCntValue = regularQueryInvocationCnt.getAndSet(0); int regularQueryInvocationCntValue = regularQueryInvocationCnt.getAndSet(0);
long regularQueryInvocationTimeValue = regularQueryTimeSpent.getAndSet(0); long regularQueryInvocationTimeValue = regularQueryTimeSpent.getAndSet(0);
@ -263,17 +307,25 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
} }
private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) {
Map<Integer, TbEntityDataSubCtx> sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>());
TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, sessionRef, cmd.getCmdId()); TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, sessionRef, cmd.getCmdId());
ctx.setQuery(cmd.getQuery()); ctx.setQuery(cmd.getQuery());
sessionSubs.put(cmd.getCmdId(), ctx); sessionSubs.put(cmd.getCmdId(), ctx);
return ctx; return ctx;
} }
private TbEntityDataSubCtx getSubCtx(String sessionId, int cmdId) { private TbAlarmDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) {
Map<Integer, TbEntityDataSubCtx> sessionSubs = subscriptionsBySessionId.get(sessionId); Map<Integer, TbAbstractDataSubCtx> 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 extends TbAbstractDataSubCtx> T getSubCtx(String sessionId, int cmdId) {
Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.get(sessionId);
if (sessionSubs != null) { if (sessionSubs != null) {
return sessionSubs.get(cmdId); return (T) sessionSubs.get(cmdId);
} else { } else {
return null; 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()))); return data.stream().collect(Collectors.toMap(TsKvEntry::getKey, value -> new TsValue(value.getTs(), value.getValueAsString())));
} }
private Map<String, List<TsValue>> toTsValues(List<TsKvEntry> data) {
Map<String, List<TsValue>> results = new HashMap<>();
for (TsKvEntry tsKvEntry : data) {
results.computeIfAbsent(tsKvEntry.getKey(), k -> new ArrayList<>()).add(new TsValue(tsKvEntry.getTs(), tsKvEntry.getValueAsString()));
}
return results;
}
@Override @Override
public void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd cmd) { public void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd cmd) {
cleanupAndCancel(getSubCtx(sessionId, cmd.getCmdId())); cleanupAndCancel(getSubCtx(sessionId, cmd.getCmdId()));
@ -442,9 +486,9 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
@Override @Override
public void cancelAllSessionSubscriptions(String sessionId) { public void cancelAllSessionSubscriptions(String sessionId) {
Map<Integer, TbEntityDataSubCtx> sessionSubs = subscriptionsBySessionId.remove(sessionId); Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.remove(sessionId);
if (sessionSubs != null) { if (sessionSubs != null) {
sessionSubs.values().forEach(this::cleanupAndCancel); sessionSubs.values().stream().filter(sub -> sub instanceof TbEntityDataSubCtx).map(sub -> (TbEntityDataSubCtx) sub).forEach(this::cleanupAndCancel);
} }
} }

59
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<T extends AbstractDataQuery> {
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();
}
}

63
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<AlarmDataQuery> {
@Getter
@Setter
private final LinkedHashMap<EntityId, EntityData> 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<EntityData> entitiesData) {
entitiesMap.clear();
tooManyEntities = entitiesData.hasNext();
for (EntityData entityData : entitiesData.getData()) {
entitiesMap.put(entityData.getEntityId(), entityData);
}
}
public Collection<EntityId> getOrderedEntityIds() {
return entitiesMap.keySet();
}
public void setAlarmsData(PageData<AlarmData> alarms) {
// TODO: implement
}
}

32
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.AllArgsConstructor;
import lombok.Data; import lombok.Data;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
@ -51,17 +54,13 @@ import java.util.function.Function;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@Slf4j @Slf4j
@Data public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
public class TbEntityDataSubCtx {
public static final int MAX_SUBS_PER_CMD = 1024 * 8; @Getter @Setter
private final String serviceId;
private final TelemetryWebSocketService wsService;
private final TelemetryWebSocketSessionRef sessionRef;
private final int cmdId;
private EntityDataQuery query;
private TimeSeriesCmd tsCmd; private TimeSeriesCmd tsCmd;
@Getter
private PageData<EntityData> data; private PageData<EntityData> data;
@Getter @Setter
private boolean initialDataSent; private boolean initialDataSent;
private Map<Integer, EntityId> subToEntityIdMap; private Map<Integer, EntityId> subToEntityIdMap;
private volatile ScheduledFuture<?> refreshTask; private volatile ScheduledFuture<?> refreshTask;
@ -69,22 +68,7 @@ public class TbEntityDataSubCtx {
private LatestValueCmd latestValueCmd; private LatestValueCmd latestValueCmd;
public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) {
this.serviceId = serviceId; super(serviceId, wsService, sessionRef, cmdId);
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();
} }
public void setData(PageData<EntityData> data) { public void setData(PageData<EntityData> data) {

3
application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java

@ -16,6 +16,7 @@
package org.thingsboard.server.service.subscription; package org.thingsboard.server.service.subscription;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; 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.EntityDataCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; 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, EntityDataCmd cmd);
void handleCmd(TelemetryWebSocketSessionRef sessionId, AlarmDataCmd cmd);
void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd subscriptionId); void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd subscriptionId);
void cancelAllSessionSubscriptions(String sessionId); void cancelAllSessionSubscriptions(String sessionId);

33
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.v1.TelemetryPluginCmd;
import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper;
import org.thingsboard.server.service.telemetry.cmd.v1.TimeseriesSubscriptionCmd; 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.EntityDataCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate;
@ -210,6 +213,9 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
if (cmdsWrapper.getEntityDataCmds() != null) { if (cmdsWrapper.getEntityDataCmds() != null) {
cmdsWrapper.getEntityDataCmds().forEach(cmd -> handleWsEntityDataCmd(sessionRef, cmd)); cmdsWrapper.getEntityDataCmds().forEach(cmd -> handleWsEntityDataCmd(sessionRef, cmd));
} }
if (cmdsWrapper.getAlarmDataCmds() != null) {
cmdsWrapper.getAlarmDataCmds().forEach(cmd -> handleWsAlarmDataCmd(sessionRef, cmd));
}
if (cmdsWrapper.getEntityDataUnsubscribeCmds() != null) { if (cmdsWrapper.getEntityDataUnsubscribeCmds() != null) {
cmdsWrapper.getEntityDataUnsubscribeCmds().forEach(cmd -> handleWsEntityDataUnsubscribeCmd(sessionRef, cmd)); 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) { private void handleWsEntityDataUnsubscribeCmd(TelemetryWebSocketSessionRef sessionRef, EntityDataUnsubscribeCmd cmd) {
String sessionId = sessionRef.getSessionId(); String sessionId = sessionRef.getSessionId();
log.debug("[{}] Processing: {}", sessionId, cmd); log.debug("[{}] Processing: {}", sessionId, cmd);
@ -246,7 +262,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
} }
@Override @Override
public void sendWsMsg(String sessionId, EntityDataUpdate update) { public void sendWsMsg(String sessionId, DataUpdate update) {
sendWsMsg(sessionId, update.getCmdId(), update); sendWsMsg(sessionId, update.getCmdId(), update);
} }
@ -661,6 +677,21 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
return true; 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) { private boolean validateSubscriptionCmd(TelemetryWebSocketSessionRef sessionRef, SubscriptionCmd cmd) {
if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty()) { if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty()) {
SubscriptionUpdate update = new SubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, SubscriptionUpdate update = new SubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST,

3
application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.service.telemetry; 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.cmd.v2.EntityDataUpdate;
import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; 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, SubscriptionUpdate update);
void sendWsMsg(String sessionId, EntityDataUpdate update); void sendWsMsg(String sessionId, DataUpdate update);
} }

6
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.AttributesSubscriptionCmd;
import org.thingsboard.server.service.telemetry.cmd.v1.GetHistoryCmd; 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.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.EntityDataCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd;
@ -40,4 +42,8 @@ public class TelemetryPluginCmdsWrapper {
private List<EntityDataUnsubscribeCmd> entityDataUnsubscribeCmds; private List<EntityDataUnsubscribeCmd> entityDataUnsubscribeCmds;
private List<AlarmDataCmd> alarmDataCmds;
private List<AlarmDataUnsubscribeCmd> alarmDataUnsubscribeCmds;
} }

33
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;
}
}

25
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;
}

46
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<AlarmData> {
public AlarmDataUpdate(int cmdId, PageData<AlarmData> data, List<AlarmData> 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<AlarmData> data,
@JsonProperty("update") List<AlarmData> update,
@JsonProperty("errorCode") int errorCode,
@JsonProperty("errorMsg") String errorMsg) {
super(cmdId, data, update, errorCode, errorMsg);
}
}

32
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;
}
}

43
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<T> {
private final int cmdId;
private final PageData<T> data;
private final List<T> update;
private final int errorCode;
private final String errorMsg;
public DataUpdate(int cmdId, PageData<T> data, List<T> update) {
this(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null);
}
public DataUpdate(int cmdId, int errorCode, String errorMsg) {
this(cmdId, null, null, errorCode, errorMsg);
}
}

24
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; 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; import org.thingsboard.server.common.data.query.EntityDataQuery;
@Data public class EntityDataCmd extends DataCmd {
public class EntityDataCmd {
private final int cmdId; @Getter
private final EntityDataQuery query; private final EntityDataQuery query;
@Getter
private final EntityHistoryCmd historyCmd; private final EntityHistoryCmd historyCmd;
@Getter
private final LatestValueCmd latestCmd; private final LatestValueCmd latestCmd;
@Getter
private final TimeSeriesCmd tsCmd; 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;
}
} }

27
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; package org.thingsboard.server.service.telemetry.cmd.v2;
import lombok.AllArgsConstructor; import com.fasterxml.jackson.annotation.JsonCreator;
import lombok.Data; import com.fasterxml.jackson.annotation.JsonProperty;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityData;
import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode;
import java.util.List; import java.util.List;
@Data public class EntityDataUpdate extends DataUpdate<EntityData> {
@AllArgsConstructor
public class EntityDataUpdate {
private final int cmdId;
private final PageData<EntityData> data;
private final List<EntityData> update;
private final int errorCode;
private final String errorMsg;
public EntityDataUpdate(int cmdId, PageData<EntityData> data, List<EntityData> update) { public EntityDataUpdate(int cmdId, PageData<EntityData> data, List<EntityData> 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) { 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<EntityData> data,
@JsonProperty("update") List<EntityData> update,
@JsonProperty("errorCode") int errorCode,
@JsonProperty("errorMsg") String errorMsg) {
super(cmdId, data, update, errorCode, errorMsg);
} }
} }

7
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_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_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}" 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:
dynamic_page_link_refresh_pool_size: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_POOL_SIZE:1}" 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: rest:
limits: limits:
tenant: tenant:

4
application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java

@ -26,9 +26,9 @@ import java.util.Arrays;
@RunWith(ClasspathSuite.class) @RunWith(ClasspathSuite.class)
@ClasspathSuite.ClassnameFilters({ @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.EntityQueryControllerSqlTest",
"org.thingsboard.server.controller.sql.*Test", // "org.thingsboard.server.controller.sql.*Test",
}) })
public class ControllerSqlTestSuite { public class ControllerSqlTestSuite {

8
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.AlarmSearchStatus;
import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmStatus; 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.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; 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. * Created by ashvayka on 11.05.17.
@ -52,4 +58,6 @@ public interface AlarmService {
ListenableFuture<Alarm> findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type); ListenableFuture<Alarm> findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type);
PageData<AlarmData> findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId,
AlarmDataPageLink pageLink, Collection<EntityId> orderedEntityIds);
} }

56
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<T extends EntityDataPageLink> extends EntityCountQuery {
@Getter
protected T pageLink;
@Getter
protected List<EntityKey> entityFields;
@Getter
protected List<EntityKey> latestValues;
@Getter
protected List<KeyFilter> keyFilters;
public AbstractDataQuery() {
super();
}
public AbstractDataQuery(EntityFilter entityFilter) {
super(entityFilter);
}
public AbstractDataQuery(EntityFilter entityFilter,
T pageLink,
List<EntityKey> entityFields,
List<EntityKey> latestValues,
List<KeyFilter> keyFilters) {
super(entityFilter);
this.pageLink = pageLink;
this.entityFields = entityFields;
this.latestValues = latestValues;
this.keyFilters = keyFilters;
}
}

37
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<EntityKeyType, Map<String, TsValue>> latest;
public AlarmData(Alarm alarm, String originatorName, UUID entityId) {
super(alarm, originatorName);
this.entityId = entityId;
this.latest = new HashMap<>();
}
}

67
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<String> typeList;
private List<AlarmSearchStatus> statusList;
private List<AlarmSeverity> 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<String> typeList, List<AlarmSearchStatus> statusList, List<AlarmSeverity> 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
);
}
}

42
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<AlarmDataPageLink> {
public AlarmDataQuery() {
}
public AlarmDataQuery(EntityFilter entityFilter) {
super(entityFilter);
}
public AlarmDataQuery(EntityFilter entityFilter, AlarmDataPageLink pageLink, List<EntityKey> entityFields, List<EntityKey> latestValues, List<KeyFilter> keyFilters) {
super(entityFilter, pageLink, entityFields, latestValues, keyFilters);
}
@JsonIgnore
public AlarmDataQuery next() {
return new AlarmDataQuery(getEntityFilter(), getPageLink().nextPageLink(), entityFields, latestValues, keyFilters);
}
}

2
common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataPageLink.java

@ -38,6 +38,6 @@ public class EntityDataPageLink {
@JsonIgnore @JsonIgnore
public EntityDataPageLink nextPageLink() { 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);
} }
} }

27
common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java

@ -22,39 +22,22 @@ import lombok.ToString;
import java.util.List; import java.util.List;
@ToString @ToString
public class EntityDataQuery extends EntityCountQuery { public class EntityDataQuery extends AbstractDataQuery<EntityDataPageLink> {
@Getter
private EntityDataPageLink pageLink;
@Getter
private List<EntityKey> entityFields;
@Getter
private List<EntityKey> latestValues;
@Getter
private List<KeyFilter> keyFilters;
public EntityDataQuery() { public EntityDataQuery() {
super();
} }
public EntityDataQuery(EntityFilter entityFilter) { public EntityDataQuery(EntityFilter entityFilter) {
super(entityFilter); super(entityFilter);
} }
public EntityDataQuery(EntityFilter entityFilter, public EntityDataQuery(EntityFilter entityFilter, EntityDataPageLink pageLink, List<EntityKey> entityFields, List<EntityKey> latestValues, List<KeyFilter> keyFilters) {
EntityDataPageLink pageLink, super(entityFilter, pageLink, entityFields, latestValues, keyFilters);
List<EntityKey> entityFields,
List<EntityKey> latestValues,
List<KeyFilter> keyFilters) {
super(entityFilter);
this.pageLink = pageLink;
this.entityFields = entityFields;
this.latestValues = latestValues;
this.keyFilters = keyFilters;
} }
@JsonIgnore @JsonIgnore
public EntityDataQuery next() { public EntityDataQuery next() {
return new EntityDataQuery(getEntityFilter(), pageLink.nextPageLink(), entityFields, latestValues, keyFilters); return new EntityDataQuery(getEntityFilter(), getPageLink().nextPageLink(), entityFields, latestValues, keyFilters);
} }
} }

3
common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKeyType.java

@ -21,5 +21,6 @@ public enum EntityKeyType {
SHARED_ATTRIBUTE, SHARED_ATTRIBUTE,
SERVER_ATTRIBUTE, SERVER_ATTRIBUTE,
TIME_SERIES, TIME_SERIES,
ENTITY_FIELD; ENTITY_FIELD,
ALARM_FIELD;
} }

8
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.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmQuery; 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.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; 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 org.thingsboard.server.dao.Dao;
import java.util.Collection;
import java.util.UUID; import java.util.UUID;
/** /**
@ -40,4 +45,7 @@ public interface AlarmDao extends Dao<Alarm> {
Alarm save(TenantId tenantId, Alarm alarm); Alarm save(TenantId tenantId, Alarm alarm);
PageData<AlarmInfo> findAlarms(TenantId tenantId, AlarmQuery query); PageData<AlarmInfo> findAlarms(TenantId tenantId, AlarmQuery query);
PageData<AlarmData> findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId,
AlarmDataPageLink pageLink, Collection<EntityId> orderedEntityIds);
} }

15
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.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.id.AlarmId; 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.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink; 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.EntityRelation;
import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.EntityRelationsQuery;
import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.EntitySearchDirection;
@ -54,6 +58,7 @@ import javax.annotation.Nullable;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collection;
import java.util.Comparator; import java.util.Comparator;
import java.util.List; import java.util.List;
import java.util.Set; import java.util.Set;
@ -69,6 +74,8 @@ import static org.thingsboard.server.dao.service.Validator.validateId;
@Slf4j @Slf4j
public class BaseAlarmService extends AbstractEntityService implements AlarmService { 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_"; public static final String ALARM_RELATION_PREFIX = "ALARM_";
@Autowired @Autowired
@ -123,6 +130,14 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
return alarmDao.findLatestByOriginatorAndType(tenantId, originator, type); return alarmDao.findLatestByOriginatorAndType(tenantId, originator, type);
} }
@Override
public PageData<AlarmData> findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId,
AlarmDataPageLink pageLink, Collection<EntityId> orderedEntityIds) {
validateId(tenantId, INCORRECT_TENANT_ID + tenantId);
validateId(customerId, INCORRECT_CUSTOMER_ID + customerId);
return alarmDao.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds);
}
@Override @Override
public Boolean deleteAlarm(TenantId tenantId, AlarmId alarmId) { public Boolean deleteAlarm(TenantId tenantId, AlarmId alarmId) {
try { try {

1
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_TYPE_PROPERTY = "type";
public static final String ALARM_DETAILS_PROPERTY = "details"; public static final String ALARM_DETAILS_PROPERTY = "details";
public static final String ALARM_ORIGINATOR_ID_PROPERTY = "originator_id"; 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_ORIGINATOR_TYPE_PROPERTY = "originator_type";
public static final String ALARM_SEVERITY_PROPERTY = "severity"; public static final String ALARM_SEVERITY_PROPERTY = "severity";
public static final String ALARM_STATUS_PROPERTY = "status"; public static final String ALARM_STATUS_PROPERTY = "status";

9
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.jpa.repository.Query;
import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.CrudRepository;
import org.springframework.data.repository.query.Param; 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.AlarmEntity;
import org.thingsboard.server.dao.model.sql.AlarmInfoEntity; import org.thingsboard.server.dao.model.sql.AlarmInfoEntity;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.util.Collection;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
@ -72,4 +79,6 @@ public interface AlarmRepository extends CrudRepository<AlarmEntity, UUID> {
@Param("endTime") Long endTime, @Param("endTime") Long endTime,
@Param("searchText") String searchText, @Param("searchText") String searchText,
Pageable pageable); Pageable pageable);
} }

14
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.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmQuery; import org.thingsboard.server.common.data.alarm.AlarmQuery;
import org.thingsboard.server.common.data.alarm.AlarmSearchStatus; 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.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; 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.DaoUtil;
import org.thingsboard.server.dao.alarm.AlarmDao; import org.thingsboard.server.dao.alarm.AlarmDao;
import org.thingsboard.server.dao.alarm.BaseAlarmService; import org.thingsboard.server.dao.alarm.BaseAlarmService;
import org.thingsboard.server.dao.model.sql.AlarmEntity; import org.thingsboard.server.dao.model.sql.AlarmEntity;
import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.relation.RelationDao;
import org.thingsboard.server.dao.sql.JpaAbstractDao; import org.thingsboard.server.dao.sql.JpaAbstractDao;
import org.thingsboard.server.dao.sql.query.AlarmQueryRepository;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.util.Collection;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
import java.util.UUID; import java.util.UUID;
@ -55,6 +61,9 @@ public class JpaAlarmDao extends JpaAbstractDao<AlarmEntity, Alarm> implements A
@Autowired @Autowired
private AlarmRepository alarmRepository; private AlarmRepository alarmRepository;
@Autowired
private AlarmQueryRepository alarmQueryRepository;
@Autowired @Autowired
private RelationDao relationDao; private RelationDao relationDao;
@ -116,4 +125,9 @@ public class JpaAlarmDao extends JpaAbstractDao<AlarmEntity, Alarm> implements A
) )
); );
} }
@Override
public PageData<AlarmData> findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataPageLink pageLink, Collection<EntityId> orderedEntityIds) {
return alarmQueryRepository.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds);
}
} }

100
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<AlarmData> createAlarmData(EntityDataPageLink pageLink,
List<Map<String, Object>> 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<AlarmData> entitiesData = convertListToAlarmData(rows);
return new PageData<>(entitiesData, totalPages, totalElements, hasNext);
}
private static List<AlarmData> convertListToAlarmData(List<Map<String, Object>> result) {
return result.stream().map(AlarmDataAdapter::toEntityData).collect(Collectors.toList());
}
private static AlarmData toEntityData(Map<String, Object> 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);
}
}

32
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<AlarmData> findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId,
AlarmDataPageLink pageLink, Collection<EntityId> orderedEntityIds);
}

241
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<String, String> 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<AlarmData> findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId,
AlarmDataPageLink pageLink, Collection<EntityId> 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<AlarmStatus> 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<Map<String, Object>> rows = jdbcTemplate.queryForList(dataQuery, ctx);
return AlarmDataAdapter.createAlarmData(pageLink, rows, totalElements);
}
private Set<AlarmStatus> toStatusSet(List<AlarmSearchStatus> statusList) {
Set<AlarmStatus> 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 ");
}
}
}

30
dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java

@ -218,7 +218,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
@Override @Override
public long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query) { public long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query) {
EntityType entityType = resolveEntityType(query.getEntityFilter()); EntityType entityType = resolveEntityType(query.getEntityFilter());
EntityQueryContext ctx = new EntityQueryContext(); QueryContext ctx = new QueryContext();
ctx.append("select count(e.id) from "); ctx.append("select count(e.id) from ");
ctx.append(addEntityTableQuery(ctx, query.getEntityFilter(), entityType)); ctx.append(addEntityTableQuery(ctx, query.getEntityFilter(), entityType));
ctx.append(" e where "); ctx.append(" e where ");
@ -228,7 +228,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
@Override @Override
public PageData<EntityData> findEntityDataByQuery(TenantId tenantId, CustomerId customerId, EntityDataQuery query) { public PageData<EntityData> findEntityDataByQuery(TenantId tenantId, CustomerId customerId, EntityDataQuery query) {
EntityQueryContext ctx = new EntityQueryContext(); QueryContext ctx = new QueryContext();
EntityType entityType = resolveEntityType(query.getEntityFilter()); EntityType entityType = resolveEntityType(query.getEntityFilter());
EntityDataPageLink pageLink = query.getPageLink(); EntityDataPageLink pageLink = query.getPageLink();
@ -308,7 +308,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
return EntityDataAdapter.createEntityData(pageLink, selectionMapping, rows, totalElements); return EntityDataAdapter.createEntityData(pageLink, selectionMapping, rows, totalElements);
} }
private String buildEntityWhere(EntityQueryContext ctx, private String buildEntityWhere(QueryContext ctx,
TenantId tenantId, TenantId tenantId,
CustomerId customerId, CustomerId customerId,
EntityFilter entityFilter, EntityFilter entityFilter,
@ -327,7 +327,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
return result; 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()) { switch (entityFilter.getType()) {
case RELATIONS_QUERY: case RELATIONS_QUERY:
case DEVICE_SEARCH_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()); ctx.addUuidParameter("permissions_tenant_id", tenantId.getId());
if (customerId != null && !customerId.isNullUid()) { if (customerId != null && !customerId.isNullUid()) {
ctx.addUuidParameter("permissions_customer_id", customerId.getId()); 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()) { switch (entityFilter.getType()) {
case SINGLE_ENTITY: case SINGLE_ENTITY:
return this.singleEntityQuery(ctx, (SingleEntityFilter) entityFilter); 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()) { switch (entityFilter.getType()) {
case RELATIONS_QUERY: case RELATIONS_QUERY:
return relationQuery(ctx, (RelationsQueryFilter) entityFilter); return relationQuery(ctx, (RelationsQueryFilter) entityFilter);
@ -393,7 +393,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
} }
} }
private String entitySearchQuery(EntityQueryContext ctx, EntitySearchQueryFilter entityFilter, EntityType entityType, List<String> types) { private String entitySearchQuery(QueryContext ctx, EntitySearchQueryFilter entityFilter, EntityType entityType, List<String> types) {
EntityId rootId = entityFilter.getRootEntity(); EntityId rootId = entityFilter.getRootEntity();
//TODO: fetch last level only. //TODO: fetch last level only.
//TODO: fetch distinct records. //TODO: fetch distinct records.
@ -416,7 +416,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
return query; return query;
} }
private String relationQuery(EntityQueryContext ctx, RelationsQueryFilter entityFilter) { private String relationQuery(QueryContext ctx, RelationsQueryFilter entityFilter) {
EntityId rootId = entityFilter.getRootEntity(); EntityId rootId = entityFilter.getRootEntity();
String lvlFilter = getLvlFilter(entityFilter.getMaxLevel()); String lvlFilter = getLvlFilter(entityFilter.getMaxLevel());
String selectFields = SELECT_TENANT_ID + ", " + SELECT_CUSTOMER_ID String selectFields = SELECT_TENANT_ID + ", " + SELECT_CUSTOMER_ID
@ -478,7 +478,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
return from; return from;
} }
private String buildWhere(EntityQueryContext ctx, List<EntityKeyMapping> latestFiltersMapping) { private String buildWhere(QueryContext ctx, List<EntityKeyMapping> latestFiltersMapping) {
String latestFilters = EntityKeyMapping.buildQuery(ctx, latestFiltersMapping); String latestFilters = EntityKeyMapping.buildQuery(ctx, latestFiltersMapping);
if (!StringUtils.isEmpty(latestFilters)) { if (!StringUtils.isEmpty(latestFilters)) {
return String.format("where %s", latestFilters); return String.format("where %s", latestFilters);
@ -487,7 +487,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
} }
} }
private String buildTextSearchQuery(EntityQueryContext ctx, List<EntityKeyMapping> selectionMapping, String searchText) { private String buildTextSearchQuery(QueryContext ctx, List<EntityKeyMapping> selectionMapping, String searchText) {
if (!StringUtils.isEmpty(searchText) && !selectionMapping.isEmpty()) { if (!StringUtils.isEmpty(searchText) && !selectionMapping.isEmpty()) {
String lowerSearchText = searchText.toLowerCase() + "%"; String lowerSearchText = searchText.toLowerCase() + "%";
List<String> searchPredicates = selectionMapping.stream().map(mapping -> { List<String> 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()); ctx.addUuidParameter("entity_filter_single_entity_id", filter.getSingleEntity().getId());
return "e.id=:entity_filter_single_entity_id"; 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())); ctx.addUuidListParameter("entity_filter_entity_ids", filter.getEntityList().stream().map(UUID::fromString).collect(Collectors.toList()));
return "e.id in (:entity_filter_entity_ids)"; 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()); ctx.addStringParameter("entity_filter_name_filter", filter.getEntityNameFilter());
return "lower(e.search_text) like lower(concat(:entity_filter_name_filter, '%%'))"; 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 type;
String name; String name;
switch (filter.getType()) { switch (filter.getType()) {

22
dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java

@ -194,7 +194,7 @@ public class EntityKeyMapping {
} }
} }
public Stream<String> toQueries(EntityQueryContext ctx) { public Stream<String> toQueries(QueryContext ctx) {
if (hasFilter()) { if (hasFilter()) {
String keyAlias = entityKey.getType().equals(EntityKeyType.ENTITY_FIELD) ? "e" : alias; String keyAlias = entityKey.getType().equals(EntityKeyType.ENTITY_FIELD) ? "e" : alias;
return keyFilters.stream().map(keyFilter -> 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; String entityTypeStr;
if (entityFilter.getType().equals(EntityFilterType.RELATIONS_QUERY)) { if (entityFilter.getType().equals(EntityFilterType.RELATIONS_QUERY)) {
entityTypeStr = "entities.entity_type"; entityTypeStr = "entities.entity_type";
@ -239,12 +239,12 @@ public class EntityKeyMapping {
Collectors.joining(", ")); Collectors.joining(", "));
} }
public static String buildLatestJoins(EntityQueryContext ctx, EntityFilter entityFilter, EntityType entityType, List<EntityKeyMapping> latestMappings) { public static String buildLatestJoins(QueryContext ctx, EntityFilter entityFilter, EntityType entityType, List<EntityKeyMapping> latestMappings) {
return latestMappings.stream().map(mapping -> mapping.toLatestJoin(ctx, entityFilter, entityType)).collect( return latestMappings.stream().map(mapping -> mapping.toLatestJoin(ctx, entityFilter, entityType)).collect(
Collectors.joining(" ")); Collectors.joining(" "));
} }
public static String buildQuery(EntityQueryContext ctx, List<EntityKeyMapping> mappings) { public static String buildQuery(QueryContext ctx, List<EntityKeyMapping> mappings) {
return mappings.stream().flatMap(mapping -> mapping.toQueries(ctx)).collect( return mappings.stream().flatMap(mapping -> mapping.toQueries(ctx)).collect(
Collectors.joining(" AND ")); Collectors.joining(" AND "));
} }
@ -357,11 +357,11 @@ public class EntityKeyMapping {
return String.join(", ", attrValSelection, attrTsSelection); 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()); 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)) { if (predicate.getType().equals(FilterPredicateType.COMPLEX)) {
return this.buildComplexPredicateQuery(ctx, alias, key, (ComplexFilterPredicate) predicate); return this.buildComplexPredicateQuery(ctx, alias, key, (ComplexFilterPredicate) predicate);
} else { } 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() return predicate.getPredicates().stream()
.map(keyFilterPredicate -> this.buildPredicateQuery(ctx, alias, key, keyFilterPredicate)).collect(Collectors.joining( .map(keyFilterPredicate -> this.buildPredicateQuery(ctx, alias, key, keyFilterPredicate)).collect(Collectors.joining(
" " + predicate.getOperation().name() + " " " " + 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 (predicate.getType().equals(FilterPredicateType.NUMERIC)) {
if (key.getType().equals(EntityKeyType.ENTITY_FIELD)) { if (key.getType().equals(EntityKeyType.ENTITY_FIELD)) {
String column = entityFieldColumnMap.get(key.getKey()); 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 operationField = field;
String paramName = getNextParameterName(field); String paramName = getNextParameterName(field);
String value = stringFilterPredicate.getValue(); String value = stringFilterPredicate.getValue();
@ -439,7 +439,7 @@ public class EntityKeyMapping {
return String.format("(%s is not null and %s)", field, stringOperationQuery); 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); String paramName = getNextParameterName(field);
ctx.addDoubleParameter(paramName, numericFilterPredicate.getValue()); ctx.addDoubleParameter(paramName, numericFilterPredicate.getValue());
String numericOperationQuery = ""; String numericOperationQuery = "";
@ -466,7 +466,7 @@ public class EntityKeyMapping {
return String.format("(%s is not null and %s)", field, numericOperationQuery); 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) { BooleanFilterPredicate booleanFilterPredicate) {
String paramName = getNextParameterName(field); String paramName = getNextParameterName(field);
ctx.addBooleanParameter(paramName, booleanFilterPredicate.isValue()); ctx.addBooleanParameter(paramName, booleanFilterPredicate.isValue());

8
dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityQueryContext.java → 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.Map;
import java.util.UUID; import java.util.UUID;
public class EntityQueryContext implements SqlParameterSource { public class QueryContext implements SqlParameterSource {
private static final PostgresUUIDType UUID_TYPE = new PostgresUUIDType(); private static final PostgresUUIDType UUID_TYPE = new PostgresUUIDType();
private final StringBuilder query; private final StringBuilder query;
private final Map<String, Parameter> params; private final Map<String, Parameter> params;
public EntityQueryContext() { public QueryContext() {
query = new StringBuilder(); query = new StringBuilder();
params = new HashMap<>(); params = new HashMap<>();
} }
@ -91,6 +91,10 @@ public class EntityQueryContext implements SqlParameterSource {
addParameter(name, value, Types.DOUBLE, "DOUBLE"); 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<String> value) { public void addStringListParameter(String name, List<String> value) {
addParameter(name, value, Types.VARCHAR, "VARCHAR"); addParameter(name, value, Types.VARCHAR, "VARCHAR");
} }

132
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.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmQuery; 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.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.id.AssetId; 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.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.SortOrder;
import org.thingsboard.server.common.data.page.TimePageLink; 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.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import java.util.Arrays;
import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
@ -195,6 +205,128 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest {
Assert.assertEquals(created, alarms.getData().get(0)); 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<AlarmData> 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 @Test
public void testDeleteAlarm() throws ExecutionException, InterruptedException { public void testDeleteAlarm() throws ExecutionException, InterruptedException {
AssetId parentId = new AssetId(Uuids.timeBased()); AssetId parentId = new AssetId(Uuids.timeBased());

4
ui-ngx/package-lock.json

@ -8997,10 +8997,10 @@
"integrity": "sha512-4O3GWAYJaauMCILm07weko2rHA8a4kjn7+8Lg4s1d7SxwS/3IpkVD/GljbRrIJ1c1W/XGJ3GbuK7RyYZEJChhw==" "integrity": "sha512-4O3GWAYJaauMCILm07weko2rHA8a4kjn7+8Lg4s1d7SxwS/3IpkVD/GljbRrIJ1c1W/XGJ3GbuK7RyYZEJChhw=="
}, },
"ngx-flowchart": { "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", "from": "git://github.com/thingsboard/ngx-flowchart.git#master",
"requires": { "requires": {
"tslib": "^1.13.0" "tslib": "^1.10.0"
}, },
"dependencies": { "dependencies": {
"tslib": { "tslib": {

Loading…
Cancel
Save