From fa8535ca95346ee94e82db9238822315d18f1f09 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 9 Jul 2020 13:23:46 +0300 Subject: [PATCH] Alarm Search Text --- ...efaultTbEntityDataSubscriptionService.java | 35 ++--- .../subscription/TbAbstractDataSubCtx.java | 121 ++++++++++++++++- .../subscription/TbAlarmDataSubCtx.java | 52 +++++++- .../subscription/TbEntityDataSubCtx.java | 126 +----------------- .../DefaultAlarmSubscriptionService.java | 5 +- .../server/dao/alarm/AlarmService.java | 2 +- .../common/data/query/AlarmDataQuery.java | 9 +- .../server/dao/alarm/AlarmDao.java | 2 +- .../server/dao/alarm/BaseAlarmService.java | 4 +- .../server/dao/sql/alarm/AlarmRepository.java | 6 +- .../server/dao/sql/alarm/JpaAlarmDao.java | 8 +- .../dao/sql/audit/AuditLogRepository.java | 8 +- .../dao/sql/query/AlarmQueryRepository.java | 3 +- .../query/DefaultAlarmQueryRepository.java | 40 +++++- .../dao/service/BaseAlarmServiceTest.java | 26 ++-- .../engine/api/RuleEngineAlarmService.java | 3 +- 16 files changed, 263 insertions(+), 187 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index 5b83fc5212..f1616e5caf 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -164,7 +164,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc if (ctx != null) { log.debug("[{}][{}] Updating existing subscriptions using: {}", session.getSessionId(), cmd.getCmdId(), cmd); if (cmd.getLatestCmd() != null || cmd.getTsCmd() != null || cmd.getHistoryCmd() != null) { - clearSubs(ctx); + ctx.clearSubscriptions(); } } else { log.debug("[{}][{}] Creating new subscription using: {}", session.getSessionId(), cmd.getCmdId(), cmd); @@ -259,19 +259,19 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc EntityDataQuery edq = new EntityDataQuery(adq.getEntityFilter(), edpl, adq.getEntityFields(), adq.getLatestValues(), adq.getKeyFilters()); PageData entitiesData = entityService.findEntityDataByQuery(ctx.getTenantId(), ctx.getCustomerId(), edq); List entities = entitiesData.getData(); - ctx.setEntitiesData(entitiesData); + ctx.setData(entitiesData); ctx.cancelTasks(); + ctx.clearSubscriptions(); if (entities.isEmpty()) { AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(Collections.emptyList(), 1, 0, false), null); wsService.sendWsMsg(ctx.getSessionId(), update); } else { ctx.fetchAlarms(); + ctx.createSubscriptions(cmd.getQuery().getLatestValues(), true); if (adq.getPageLink().getTimeWindow() > 0) { - ctx.createSubscriptions(); TbAlarmDataSubCtx finalCtx = ctx; ScheduledFuture task = scheduler.scheduleWithFixedDelay( - finalCtx::cleanupOldAlarms, - dynamicPageLinkRefreshInterval, dynamicPageLinkRefreshInterval, TimeUnit.SECONDS); + finalCtx::cleanupOldAlarms, dynamicPageLinkRefreshInterval, dynamicPageLinkRefreshInterval, TimeUnit.SECONDS); finalCtx.setRefreshTask(task); } } @@ -280,12 +280,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private void refreshDynamicQuery(TenantId tenantId, CustomerId customerId, TbEntityDataSubCtx finalCtx) { try { long start = System.currentTimeMillis(); - TbEntityDataSubCtx.TbEntityDataSubCtxUpdateResult result = finalCtx.update(entityService.findEntityDataByQuery(tenantId, customerId, finalCtx.getQuery())); + finalCtx.update(entityService.findEntityDataByQuery(tenantId, customerId, finalCtx.getQuery())); long end = System.currentTimeMillis(); dynamicQueryInvocationCnt.incrementAndGet(); dynamicQueryTimeSpent.addAndGet(end - start); - result.getSubsToCancel().forEach(subId -> localSubscriptionService.cancelSubscription(finalCtx.getSessionId(), subId)); - result.getSubsToAdd().forEach(localSubscriptionService::addSubscription); } catch (Exception e) { log.warn("[{}][{}] Failed to refresh query", finalCtx.getSessionId(), finalCtx.getCmdId(), e); } @@ -304,10 +302,6 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } - private void clearSubs(TbAbstractDataSubCtx ctx) { - ctx.clearSubscriptions(); - } - private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, localSubscriptionService, sessionRef, cmd.getCmdId()); @@ -391,7 +385,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } wsService.sendWsMsg(ctx.getSessionId(), update); if (subscribe) { - createTelemetrySubscriptions(ctx, keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), false); + ctx.createSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), false); } ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); return ctx; @@ -440,7 +434,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); } wsService.sendWsMsg(ctx.getSessionId(), update); - createTelemetrySubscriptions(ctx, latestCmd.getKeys()); + ctx.createSubscriptions(latestCmd.getKeys(), true); } @Override @@ -456,19 +450,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc wsService.sendWsMsg(ctx.getSessionId(), update); ctx.setInitialDataSent(true); } - createTelemetrySubscriptions(ctx, latestCmd.getKeys()); + ctx.createSubscriptions(latestCmd.getKeys(), true); } } - private void createTelemetrySubscriptions(TbEntityDataSubCtx ctx, List keys) { - createTelemetrySubscriptions(ctx, keys, true); - } - - private void createTelemetrySubscriptions(TbEntityDataSubCtx ctx, List keys, boolean latest) { - List tbSubs = ctx.createSubscriptions(keys, latest); - tbSubs.forEach(sub -> localSubscriptionService.addSubscription(sub)); - } - private Map toTsValue(List data) { return data.stream().collect(Collectors.toMap(TsKvEntry::getKey, value -> new TsValue(value.getTs(), value.getValueAsString()))); } @@ -481,7 +466,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private void cleanupAndCancel(TbAbstractDataSubCtx ctx) { if (ctx != null) { ctx.cancelTasks(); - clearSubs(ctx); + ctx.clearSubscriptions(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java index 244613a984..b42e23bdea 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java @@ -22,11 +22,21 @@ import lombok.extern.slf4j.Slf4j; 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.AbstractDataQuery; +import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; +import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; +import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; +import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.ScheduledFuture; @@ -39,10 +49,13 @@ public abstract class TbAbstractDataSubCtx { protected final TbLocalSubscriptionService localSubscriptionService; protected final TelemetryWebSocketSessionRef sessionRef; protected final int cmdId; + protected final Map subToEntityIdMap; + + @Getter + protected PageData data; @Getter @Setter protected T query; - protected Map subToEntityIdMap; @Setter protected volatile ScheduledFuture refreshTask; @@ -53,6 +66,7 @@ public abstract class TbAbstractDataSubCtx { this.localSubscriptionService = localSubscriptionService; this.sessionRef = sessionRef; this.cmdId = cmdId; + this.subToEntityIdMap = new HashMap<>(); } public String getSessionId() { @@ -67,7 +81,7 @@ public abstract class TbAbstractDataSubCtx { return sessionRef.getSecurityCtx().getCustomerId(); } - public void clearSubscriptions(){ + public void clearSubscriptions() { if (subToEntityIdMap != null) { for (Integer subId : subToEntityIdMap.keySet()) { localSubscriptionService.cancelSubscription(sessionRef.getSessionId(), subId); @@ -76,6 +90,10 @@ public abstract class TbAbstractDataSubCtx { } } + public void setData(PageData data) { + this.data = data; + } + public void setRefreshTask(ScheduledFuture task) { this.refreshTask = task; } @@ -87,4 +105,103 @@ public abstract class TbAbstractDataSubCtx { } } + + public void createSubscriptions(List keys, boolean resultToLatestValues) { + Map> keysByType = getEntityKeyByTypeMap(keys); + for (EntityData entityData : data.getData()) { + List entitySubscriptions = addSubscriptions(entityData, keysByType, resultToLatestValues); + entitySubscriptions.forEach(localSubscriptionService::addSubscription); + } + } + + protected Map> getEntityKeyByTypeMap(List keys) { + Map> keysByType = new HashMap<>(); + keys.forEach(key -> keysByType.computeIfAbsent(key.getType(), k -> new ArrayList<>()).add(key)); + return keysByType; + } + + protected List addSubscriptions(EntityData entityData, Map> keysByType, boolean resultToLatestValues) { + List subscriptionList = new ArrayList<>(); + keysByType.forEach((keysType, keysList) -> { + int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); + subToEntityIdMap.put(subIdx, entityData.getEntityId()); + switch (keysType) { + case TIME_SERIES: + subscriptionList.add(createTsSub(entityData, subIdx, keysList, resultToLatestValues)); + break; + case CLIENT_ATTRIBUTE: + subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.CLIENT_SCOPE, keysList)); + break; + case SHARED_ATTRIBUTE: + subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.SHARED_SCOPE, keysList)); + break; + case SERVER_ATTRIBUTE: + subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.SERVER_SCOPE, keysList)); + break; + case ATTRIBUTE: + subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.ANY_SCOPE, keysList)); + break; + } + }); + return subscriptionList; + } + + private TbSubscription createAttrSub(EntityData entityData, int subIdx, EntityKeyType keysType, TbAttributeSubscriptionScope scope, List subKeys) { + Map keyStates = buildKeyStats(entityData, keysType, subKeys); + log.trace("[{}][{}][{}] Creating attributes subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates); + return TbAttributeSubscription.builder() + .serviceId(serviceId) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(subIdx) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityData.getEntityId()) + .updateConsumer((s, subscriptionUpdate) -> sendWsMsg(s, subscriptionUpdate, keysType)) + .allKeys(false) + .keyStates(keyStates) + .scope(scope) + .build(); + } + + private TbSubscription createTsSub(EntityData entityData, int subIdx, List subKeys, boolean resultToLatestValues) { + Map keyStates = buildKeyStats(entityData, EntityKeyType.TIME_SERIES, subKeys); + if (entityData.getTimeseries() != null) { + entityData.getTimeseries().forEach((k, v) -> { + long ts = Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L); + log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, ts); + keyStates.put(k, ts); + }); + } + log.trace("[{}][{}][{}] Creating time-series subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates); + return TbTimeseriesSubscription.builder() + .serviceId(serviceId) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(subIdx) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityData.getEntityId()) + .updateConsumer((sessionId, subscriptionUpdate) -> sendWsMsg(sessionId, subscriptionUpdate, EntityKeyType.TIME_SERIES, resultToLatestValues)) + .allKeys(false) + .keyStates(keyStates) + .build(); + } + + private void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType) { + sendWsMsg(sessionId, subscriptionUpdate, keyType, true); + } + + private Map buildKeyStats(EntityData entityData, EntityKeyType keysType, List subKeys) { + Map keyStates = new HashMap<>(); + subKeys.forEach(key -> keyStates.put(key.getKey(), 0L)); + if (entityData.getLatest() != null) { + Map currentValues = entityData.getLatest().get(keysType); + if (currentValues != null) { + currentValues.forEach((k, v) -> { + log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, v.getTs()); + keyStates.put(k, v.getTs()); + }); + } + } + return keyStates; + } + + abstract void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java index c00450e871..db849ef5d5 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java @@ -27,16 +27,23 @@ 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.EntityData; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate; +import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; +import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; import java.util.function.Function; import java.util.stream.Collectors; @@ -50,6 +57,8 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { @Getter @Setter private final HashMap alarmsMap; + + private final List alarmSubscriptions; @Getter @Setter private PageData alarms; @@ -65,20 +74,22 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { this.alarmService = alarmService; this.entitiesMap = new LinkedHashMap<>(); this.alarmsMap = new HashMap<>(); + this.alarmSubscriptions = new ArrayList<>(); } public void fetchAlarms() { PageData alarms = alarmService.findAlarmDataByQueryForEntities(getTenantId(), getCustomerId(), - query.getPageLink(), getOrderedEntityIds()); + query, getOrderedEntityIds()); alarms = setAndMergeAlarmsData(alarms); AlarmDataUpdate update = new AlarmDataUpdate(cmdId, alarms, null); wsService.sendWsMsg(getSessionId(), update); } - public void setEntitiesData(PageData entitiesData) { + public void setData(PageData data) { + super.setData(data); entitiesMap.clear(); - tooManyEntities = entitiesData.hasNext(); - for (EntityData entityData : entitiesData.getData()) { + tooManyEntities = data.hasNext(); + for (EntityData entityData : data.getData()) { entitiesMap.put(entityData.getEntityId(), entityData); } } @@ -103,9 +114,13 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { return this.alarms; } - public void createSubscriptions() { - clearSubscriptions(); - this.subToEntityIdMap = new HashMap<>(); + @Override + public void createSubscriptions(List keys, boolean resultToLatestValues) { + super.createSubscriptions(keys, resultToLatestValues); + createAlarmSubscriptions(); + } + + public void createAlarmSubscriptions() { AlarmDataPageLink pageLink = query.getPageLink(); long startTs = System.currentTimeMillis() - pageLink.getTimeWindow(); for (EntityData entityData : entitiesMap.values()) { @@ -126,6 +141,28 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { } } + @Override + void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues) { + EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId()); + if (entityId != null) { + Map latestUpdate = new HashMap<>(); + subscriptionUpdate.getData().forEach((k, v) -> { + Object[] data = (Object[]) v.get(0); + latestUpdate.put(k, new TsValue((Long) data[0], (String) data[1])); + }); + EntityData entityData = entitiesMap.get(entityId); + entityData.getLatest().computeIfAbsent(keyType, tmp -> new HashMap<>()).putAll(latestUpdate); + log.trace("[{}][{}][{}][{}] Received subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); + List update = alarmsMap.values().stream().filter(alarm -> entityId.equals(alarm.getEntityId())).map(alarm -> { + alarm.getLatest().computeIfAbsent(keyType, tmp -> new HashMap<>()).putAll(latestUpdate); + return alarm; + }).collect(Collectors.toList()); + wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update)); + } else { + log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); + } + } + private void sendWsMsg(String sessionId, AlarmSubscriptionUpdate subscriptionUpdate) { Alarm alarm = subscriptionUpdate.getAlarm(); AlarmId alarmId = alarm.getId(); @@ -141,6 +178,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { if (onCurrentPage) { if (matchesFilter) { AlarmData updated = new AlarmData(alarm, current.getOriginatorName(), current.getEntityId()); + updated.getLatest().putAll(current.getLatest()); alarmsMap.put(alarmId, updated); wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated))); } else { diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index 80505b685b..f57de8e504 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -57,8 +57,6 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { @Setter private TimeSeriesCmd tsCmd; @Getter - private PageData data; - @Getter @Setter private boolean initialDataSent; private TimeSeriesCmd curTsCmd; @@ -70,110 +68,8 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { super(serviceId, wsService, localSubscriptionService, sessionRef, cmdId); } - public void setData(PageData data) { - this.data = data; - } - - public List createSubscriptions(List keys, boolean resultToLatestValues) { - this.subToEntityIdMap = new HashMap<>(); - List tbSubs = new ArrayList<>(); - Map> keysByType = getEntityKeyByTypeMap(keys); - for (EntityData entityData : data.getData()) { - tbSubs.addAll(addSubscriptions(entityData, keysByType, resultToLatestValues)); - } - return tbSubs; - } - - private Map> getEntityKeyByTypeMap(List keys) { - Map> keysByType = new HashMap<>(); - keys.forEach(key -> keysByType.computeIfAbsent(key.getType(), k -> new ArrayList<>()).add(key)); - return keysByType; - } - - private List addSubscriptions(EntityData entityData, Map> keysByType, boolean resultToLatestValues) { - List subscriptionList = new ArrayList<>(); - keysByType.forEach((keysType, keysList) -> { - int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); - subToEntityIdMap.put(subIdx, entityData.getEntityId()); - switch (keysType) { - case TIME_SERIES: - subscriptionList.add(createTsSub(entityData, subIdx, keysList, resultToLatestValues)); - break; - case CLIENT_ATTRIBUTE: - subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.CLIENT_SCOPE, keysList)); - break; - case SHARED_ATTRIBUTE: - subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.SHARED_SCOPE, keysList)); - break; - case SERVER_ATTRIBUTE: - subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.SERVER_SCOPE, keysList)); - break; - case ATTRIBUTE: - subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.ANY_SCOPE, keysList)); - break; - } - }); - return subscriptionList; - } - - private TbSubscription createAttrSub(EntityData entityData, int subIdx, EntityKeyType keysType, TbAttributeSubscriptionScope scope, List subKeys) { - Map keyStates = buildKeyStats(entityData, keysType, subKeys); - log.trace("[{}][{}][{}] Creating attributes subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates); - return TbAttributeSubscription.builder() - .serviceId(serviceId) - .sessionId(sessionRef.getSessionId()) - .subscriptionId(subIdx) - .tenantId(sessionRef.getSecurityCtx().getTenantId()) - .entityId(entityData.getEntityId()) - .updateConsumer((s, subscriptionUpdate) -> sendWsMsg(s, subscriptionUpdate, keysType)) - .allKeys(false) - .keyStates(keyStates) - .scope(scope) - .build(); - } - - private TbSubscription createTsSub(EntityData entityData, int subIdx, List subKeys, boolean resultToLatestValues) { - Map keyStates = buildKeyStats(entityData, EntityKeyType.TIME_SERIES, subKeys); - if (entityData.getTimeseries() != null) { - entityData.getTimeseries().forEach((k, v) -> { - long ts = Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L); - log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, ts); - keyStates.put(k, ts); - }); - } - log.trace("[{}][{}][{}] Creating time-series subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates); - return TbTimeseriesSubscription.builder() - .serviceId(serviceId) - .sessionId(sessionRef.getSessionId()) - .subscriptionId(subIdx) - .tenantId(sessionRef.getSecurityCtx().getTenantId()) - .entityId(entityData.getEntityId()) - .updateConsumer((sessionId, subscriptionUpdate) -> sendWsMsg(sessionId, subscriptionUpdate, EntityKeyType.TIME_SERIES, resultToLatestValues)) - .allKeys(false) - .keyStates(keyStates) - .build(); - } - - private Map buildKeyStats(EntityData entityData, EntityKeyType keysType, List subKeys) { - Map keyStates = new HashMap<>(); - subKeys.forEach(key -> keyStates.put(key.getKey(), 0L)); - if (entityData.getLatest() != null) { - Map currentValues = entityData.getLatest().get(keysType); - if (currentValues != null) { - currentValues.forEach((k, v) -> { - log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, v.getTs()); - keyStates.put(k, v.getTs()); - }); - } - } - return keyStates; - } - - private void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType) { - sendWsMsg(sessionId, subscriptionUpdate, keyType, true); - } - - private void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues) { + @Override + protected void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues) { EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId()); if (entityId != null) { log.trace("[{}][{}][{}][{}] Received subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); @@ -272,8 +168,8 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { private EntityData getDataForEntity(EntityId entityId) { return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null); } - - public TbEntityDataSubCtxUpdateResult update(PageData newData) { + + public void update(PageData newData) { Map oldDataMap; if (data != null && !data.getData().isEmpty()) { oldDataMap = data.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); @@ -283,7 +179,6 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { Map newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) { log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId); - return TbEntityDataSubCtxUpdateResult.EMPTY; } else { this.data = newData; List subIdsToCancel = new ArrayList<>(); @@ -322,7 +217,8 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { } } wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null)); - return new TbEntityDataSubCtxUpdateResult(subIdsToCancel, subsToAdd); + subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId)); + subsToAdd.forEach(localSubscriptionService::addSubscription); } } @@ -330,14 +226,4 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { curTsCmd = cmd.getTsCmd(); latestValueCmd = cmd.getLatestCmd(); } - - @Data - @AllArgsConstructor - public static class TbEntityDataSubCtxUpdateResult { - - private static TbEntityDataSubCtxUpdateResult EMPTY = new TbEntityDataSubCtxUpdateResult(Collections.emptyList(), Collections.emptyList()); - - private List subsToCancel; - private List subsToAdd; - } } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java index 3aa9cdc6c0..4dea992067 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java @@ -45,6 +45,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -149,8 +150,8 @@ public class DefaultAlarmSubscriptionService extends AbstractSubscriptionService } @Override - public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataPageLink pageLink, Collection orderedEntityIds) { - return alarmService.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds); + public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataQuery query, Collection orderedEntityIds) { + return alarmService.findAlarmDataByQueryForEntities(tenantId, customerId, query, orderedEntityIds); } @Override diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java index 1efab86d51..0d5f4e6d37 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java @@ -59,5 +59,5 @@ public interface AlarmService { ListenableFuture findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type); PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, - AlarmDataPageLink pageLink, Collection orderedEntityIds); + AlarmDataQuery query, Collection orderedEntityIds); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java index 6052a8cd5f..ef314416a2 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/AlarmDataQuery.java @@ -16,6 +16,7 @@ package org.thingsboard.server.common.data.query; import com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.Getter; import lombok.ToString; import java.util.List; @@ -23,6 +24,9 @@ import java.util.List; @ToString public class AlarmDataQuery extends AbstractDataQuery { + @Getter + protected List alarmFields; + public AlarmDataQuery() { } @@ -30,12 +34,13 @@ public class AlarmDataQuery extends AbstractDataQuery { super(entityFilter); } - public AlarmDataQuery(EntityFilter entityFilter, AlarmDataPageLink pageLink, List entityFields, List latestValues, List keyFilters) { + public AlarmDataQuery(EntityFilter entityFilter, AlarmDataPageLink pageLink, List entityFields, List latestValues, List keyFilters, List alarmFields) { super(entityFilter, pageLink, entityFields, latestValues, keyFilters); + this.alarmFields = alarmFields; } @JsonIgnore public AlarmDataQuery next() { - return new AlarmDataQuery(getEntityFilter(), getPageLink().nextPageLink(), entityFields, latestValues, keyFilters); + return new AlarmDataQuery(getEntityFilter(), getPageLink().nextPageLink(), entityFields, latestValues, keyFilters, alarmFields); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java index ab01ba341d..ff308c49f5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java @@ -47,5 +47,5 @@ public interface AlarmDao extends Dao { PageData findAlarms(TenantId tenantId, AlarmQuery query); PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, - AlarmDataPageLink pageLink, Collection orderedEntityIds); + AlarmDataQuery query, Collection orderedEntityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java index 85241baa51..3064754a1b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java @@ -133,10 +133,10 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ @Override public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, - AlarmDataPageLink pageLink, Collection orderedEntityIds) { + AlarmDataQuery query, Collection orderedEntityIds) { validateId(tenantId, INCORRECT_TENANT_ID + tenantId); validateId(customerId, INCORRECT_CUSTOMER_ID + customerId); - return alarmDao.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds); + return alarmDao.findAlarmDataByQueryForEntities(tenantId, customerId, query, orderedEntityIds); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java index cfe9ca6840..3c85cf2c60 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java @@ -57,7 +57,7 @@ public interface AlarmRepository extends CrudRepository { "AND (a.originatorId = :affectedEntityId or re.fromId IS NOT NULL) " + "AND (:startTime IS NULL OR a.createdTime >= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND (:alarmStatuses IS NULL OR a.status in :alarmStatuses) " + + "AND ((:alarmStatuses) IS NULL OR a.status in (:alarmStatuses)) " + "AND (LOWER(a.type) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.severity) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.status) LIKE LOWER(CONCAT(:searchText, '%')))", @@ -71,7 +71,7 @@ public interface AlarmRepository extends CrudRepository { "AND (a.originatorId = :affectedEntityId or re.fromId IS NOT NULL) " + "AND (:startTime IS NULL OR a.createdTime >= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND (:alarmStatuses IS NULL OR a.status in :alarmStatuses) " + + "AND ((:alarmStatuses) IS NULL OR a.status in (:alarmStatuses)) " + "AND (LOWER(a.type) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.severity) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.status) LIKE LOWER(CONCAT(:searchText, '%')))") @@ -80,7 +80,7 @@ public interface AlarmRepository extends CrudRepository { @Param("affectedEntityType") String affectedEntityType, @Param("startTime") Long startTime, @Param("endTime") Long endTime, - @Param("alarmStatuses") Set alarmStatuses, + @Param("alarmStatuses") List alarmStatuses, @Param("searchText") String searchText, Pageable pageable); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java index a4398e2708..2b42f81e6e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java @@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.alarm.AlarmDao; import org.thingsboard.server.dao.model.sql.AlarmEntity; @@ -40,6 +41,7 @@ import org.thingsboard.server.dao.sql.JpaAbstractDao; import org.thingsboard.server.dao.sql.query.AlarmQueryRepository; import org.thingsboard.server.dao.util.SqlDao; +import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.List; @@ -112,7 +114,7 @@ public class JpaAlarmDao extends JpaAbstractDao implements A affectedEntity.getEntityType().name(), query.getPageLink().getStartTime(), query.getPageLink().getEndTime(), - statusSet, + new ArrayList<>(statusSet), Objects.toString(query.getPageLink().getTextSearch(), ""), DaoUtil.toPageable(query.getPageLink()) ) @@ -120,7 +122,7 @@ public class JpaAlarmDao extends JpaAbstractDao implements A } @Override - public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataPageLink pageLink, Collection orderedEntityIds) { - return alarmQueryRepository.findAlarmDataByQueryForEntities(tenantId, customerId, pageLink, orderedEntityIds); + public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataQuery query, Collection orderedEntityIds) { + return alarmQueryRepository.findAlarmDataByQueryForEntities(tenantId, customerId, query, orderedEntityIds); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/audit/AuditLogRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/audit/AuditLogRepository.java index 545fcee083..b5374e318c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/audit/AuditLogRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/audit/AuditLogRepository.java @@ -33,7 +33,7 @@ public interface AuditLogRepository extends PagingAndSortingRepository= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND (:actionTypes IS NULL OR a.actionType in :actionTypes) " + + "AND ((:actionTypes) IS NULL OR a.actionType in (:actionTypes)) " + "AND (LOWER(a.entityType) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.entityName) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.userName) LIKE LOWER(CONCAT(:textSearch, '%'))" + @@ -53,7 +53,7 @@ public interface AuditLogRepository extends PagingAndSortingRepository= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND (:actionTypes IS NULL OR a.actionType in :actionTypes) " + + "AND ((:actionTypes) IS NULL OR a.actionType in (:actionTypes)) " + "AND (LOWER(a.entityName) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.userName) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.actionType) LIKE LOWER(CONCAT(:textSearch, '%'))" + @@ -73,7 +73,7 @@ public interface AuditLogRepository extends PagingAndSortingRepository= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND (:actionTypes IS NULL OR a.actionType in :actionTypes) " + + "AND ((:actionTypes) IS NULL OR a.actionType in (:actionTypes)) " + "AND (LOWER(a.entityType) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.entityName) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.userName) LIKE LOWER(CONCAT(:textSearch, '%'))" + @@ -93,7 +93,7 @@ public interface AuditLogRepository extends PagingAndSortingRepository= :startTime) " + "AND (:endTime IS NULL OR a.createdTime <= :endTime) " + - "AND (:actionTypes IS NULL OR a.actionType in :actionTypes) " + + "AND ((:actionTypes) IS NULL OR a.actionType in (:actionTypes)) " + "AND (LOWER(a.entityType) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.entityName) LIKE LOWER(CONCAT(:textSearch, '%'))" + "OR LOWER(a.actionType) LIKE LOWER(CONCAT(:textSearch, '%'))" + diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java index 25f8d98aac..767b227acc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/AlarmQueryRepository.java @@ -21,12 +21,13 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import java.util.Collection; public interface AlarmQueryRepository { PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, - AlarmDataPageLink pageLink, Collection orderedEntityIds); + AlarmDataQuery query, Collection orderedEntityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java index fe622ff6d7..6ceaa8dc04 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.sql.query; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; import org.springframework.stereotype.Repository; @@ -29,16 +30,20 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.EntityDataSortOrder; +import org.thingsboard.server.common.data.query.EntityKey; 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.ArrayList; import java.util.Collection; import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Set; import java.util.stream.Collectors; @@ -48,6 +53,7 @@ import java.util.stream.Collectors; public class DefaultAlarmQueryRepository implements AlarmQueryRepository { private static final Map alarmFieldColumnMap = new HashMap<>(); + private static final List uniqueAlarmFields = new ArrayList<>(); static { alarmFieldColumnMap.put("createdTime", ModelConstants.CREATED_TIME_PROPERTY); @@ -66,6 +72,8 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository { alarmFieldColumnMap.put("originator_id", ModelConstants.ALARM_ORIGINATOR_ID_PROPERTY); alarmFieldColumnMap.put("originator_type", ModelConstants.ALARM_ORIGINATOR_TYPE_PROPERTY); alarmFieldColumnMap.put("originator", "originator_name"); + + uniqueAlarmFields.addAll(new HashSet<>(alarmFieldColumnMap.values())); } public static final String SELECT_ORIGINATOR_NAME = " CASE" + @@ -109,7 +117,8 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository { @Override public PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, - AlarmDataPageLink pageLink, Collection orderedEntityIds) { + AlarmDataQuery query, Collection orderedEntityIds) { + AlarmDataPageLink pageLink = query.getPageLink(); QueryContext ctx = new QueryContext(); ctx.addUuidListParameter("entity_ids", orderedEntityIds.stream().map(EntityId::getId).collect(Collectors.toList())); @@ -208,10 +217,15 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository { } } - String countQuery = fromPart.toString() + wherePart.toString(); - int totalElements = jdbcTemplate.queryForObject(String.format("select count(*) %s", countQuery), ctx, Integer.class); + String textSearchQuery = buildTextSearchQuery(ctx, query.getAlarmFields(), pageLink.getTextSearch()); + String mainQuery = selectPart.toString() + fromPart.toString() + wherePart.toString(); + if (!textSearchQuery.isEmpty()) { + mainQuery = String.format("select * from (%s) a WHERE %s", mainQuery, textSearchQuery); + } + String countQuery = mainQuery; + int totalElements = jdbcTemplate.queryForObject(String.format("select count(*) from (%s) result", countQuery), ctx, Integer.class); - String dataQuery = selectPart.toString() + countQuery + sortPart; + String dataQuery = mainQuery + sortPart; int startIndex = pageLink.getPageSize() * pageLink.getPage(); if (pageLink.getPageSize() > 0) { @@ -221,6 +235,24 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository { return AlarmDataAdapter.createAlarmData(pageLink, rows, totalElements, orderedEntityIds); } + private String buildTextSearchQuery(QueryContext ctx, List selectionMapping, String searchText) { + if (!StringUtils.isEmpty(searchText) && selectionMapping != null && !selectionMapping.isEmpty()) { + String lowerSearchText = searchText.toLowerCase() + "%"; + List searchPredicates = selectionMapping.stream() + .map(mapping -> alarmFieldColumnMap.get(mapping.getKey())) + .filter(Objects::nonNull) + .map(mapping -> { + String paramName = mapping + "_lowerSearchText"; + ctx.addStringParameter(paramName, lowerSearchText); + return String.format("LOWER(cast(%s as varchar)) LIKE concat('%%', :%s, '%%')", mapping, paramName); + } + ).collect(Collectors.toList()); + return String.format("%s", String.join(" or ", searchPredicates)); + } else { + return ""; + } + } + private String buildPermissionsQuery(TenantId tenantId, CustomerId customerId, QueryContext ctx) { StringBuilder permissionsQuery = new StringBuilder(); ctx.addUuidParameter("permissions_tenant_id", tenantId.getId()); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java index 98b0100caf..15bd4eb905 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java @@ -38,6 +38,7 @@ 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.DeviceTypeFilter; import org.thingsboard.server.common.data.query.EntityDataSortOrder; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; @@ -261,10 +262,10 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); - PageData tenantAlarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Arrays.asList(tenantDevice.getId(), customerDevice.getId())); + PageData tenantAlarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Arrays.asList(tenantDevice.getId(), customerDevice.getId())); Assert.assertEquals(2, tenantAlarms.getData().size()); - PageData customerAlarms = alarmService.findAlarmDataByQueryForEntities(tenantId, customer.getId(), pageLink, Arrays.asList(tenantDevice.getId(), customerDevice.getId())); + PageData customerAlarms = alarmService.findAlarmDataByQueryForEntities(tenantId, customer.getId(), toQuery(pageLink), Arrays.asList(tenantDevice.getId(), customerDevice.getId())); Assert.assertEquals(1, customerAlarms.getData().size()); Assert.assertEquals(deviceAlarm, customerAlarms.getData().get(0)); @@ -277,7 +278,14 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(tenantAlarm, alarms.getData().get(0)); + } + + private AlarmDataQuery toQuery(AlarmDataPageLink pageLink){ + return toQuery(pageLink, Collections.EMPTY_LIST); + } + private AlarmDataQuery toQuery(AlarmDataPageLink pageLink, List alarmFields){ + return new AlarmDataQuery(new DeviceTypeFilter(), pageLink, null, null, null, alarmFields); } @Test @@ -314,7 +322,7 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { pageLink.setSeverityList(Arrays.asList(AlarmSeverity.CRITICAL, AlarmSeverity.WARNING)); pageLink.setStatusList(Arrays.asList(AlarmSearchStatus.ACTIVE)); - PageData alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(childId)); + PageData alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(childId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); @@ -330,13 +338,13 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { 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)); + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(childId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(created, new Alarm(alarms.getData().get(0))); pageLink.setSearchPropagatedAlarms(true); - alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), pageLink, Collections.singletonList(childId)); + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(childId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(created, new Alarm(alarms.getData().get(0))); @@ -357,7 +365,7 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { 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)); + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(childId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(created, alarms.getData().get(0)); @@ -373,7 +381,7 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { 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)); + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(parentId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(created, alarms.getData().get(0)); @@ -418,7 +426,7 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { 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)); + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(parentId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(created, alarms.getData().get(0)); @@ -436,7 +444,7 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest { 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)); + alarms = alarmService.findAlarmDataByQueryForEntities(tenantId, new CustomerId(CustomerId.NULL_UUID), toQuery(pageLink), Collections.singletonList(childId)); Assert.assertNotNull(alarms.getData()); Assert.assertEquals(1, alarms.getData().size()); Assert.assertEquals(created, alarms.getData().get(0)); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineAlarmService.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineAlarmService.java index 3ba6945544..719a9f40e5 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineAlarmService.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineAlarmService.java @@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataPageLink; +import org.thingsboard.server.common.data.query.AlarmDataQuery; import java.util.Collection; import java.util.List; @@ -60,5 +61,5 @@ public interface RuleEngineAlarmService { AlarmSeverity findHighestAlarmSeverity(TenantId tenantId, EntityId entityId, AlarmSearchStatus alarmSearchStatus, AlarmStatus alarmStatus); - PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataPageLink pageLink, Collection orderedEntityIds); + PageData findAlarmDataByQueryForEntities(TenantId tenantId, CustomerId customerId, AlarmDataQuery query, Collection orderedEntityIds); }