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 f1616e5caf..4c75486383 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 @@ -43,6 +43,7 @@ import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.alarm.AlarmService; +import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -102,6 +103,9 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Autowired private AlarmService alarmService; + @Autowired + private AttributesService attributesService; + @Autowired @Lazy private TbLocalSubscriptionService localSubscriptionService; @@ -164,7 +168,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) { - ctx.clearSubscriptions(); + ctx.clearEntitySubscriptions(); } } else { log.debug("[{}][{}] Creating new subscription using: {}", session.getSessionId(), cmd.getCmdId(), cmd); @@ -177,7 +181,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } else { log.debug("[{}][{}] Updating data using query: {}", session.getSessionId(), cmd.getCmdId(), cmd.getQuery()); } - ctx.setQuery(cmd.getQuery()); + ctx.setAndResolveQuery(cmd.getQuery()); TenantId tenantId = ctx.getTenantId(); CustomerId customerId = ctx.getCustomerId(); EntityDataQuery query = ctx.getQuery(); @@ -190,17 +194,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }); } long start = System.currentTimeMillis(); - PageData data = entityService.findEntityDataByQuery(tenantId, customerId, ctx.getQuery()); + ctx.fetchData(); long end = System.currentTimeMillis(); regularQueryInvocationCnt.incrementAndGet(); regularQueryTimeSpent.addAndGet(end - start); - - if (log.isTraceEnabled()) { - data.getData().forEach(ed -> { - log.trace("[{}][{}] EntityData: {}", session.getSessionId(), cmd.getCmdId(), ed); - }); - } - ctx.setData(data); ctx.cancelTasks(); if (ctx.getQuery().getPageLink().isDynamic()) { //TODO: validate number of dynamic page links against rate limits. Ignore dynamic flag if limit is reached. @@ -246,22 +243,12 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc log.debug("[{}][{}] Creating new alarm subscription using: {}", session.getSessionId(), cmd.getCmdId(), cmd); ctx = createSubCtx(session, cmd); } - ctx.setQuery(cmd.getQuery()); + ctx.setAndResolveQuery(cmd.getQuery()); AlarmDataQuery adq = ctx.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(maxEntitiesPerAlarmSubscription, 0, null, entitiesSortOrder); - EntityDataQuery edq = new EntityDataQuery(adq.getEntityFilter(), edpl, adq.getEntityFields(), adq.getLatestValues(), adq.getKeyFilters()); - PageData entitiesData = entityService.findEntityDataByQuery(ctx.getTenantId(), ctx.getCustomerId(), edq); - List entities = entitiesData.getData(); - ctx.setData(entitiesData); + ctx.fetchData(); + List entities = ctx.getEntitiesData(); ctx.cancelTasks(); - ctx.clearSubscriptions(); + ctx.clearEntitySubscriptions(); if (entities.isEmpty()) { AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(Collections.emptyList(), 1, 0, false), null); wsService.sendWsMsg(ctx.getSessionId(), update); @@ -280,7 +267,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private void refreshDynamicQuery(TenantId tenantId, CustomerId customerId, TbEntityDataSubCtx finalCtx) { try { long start = System.currentTimeMillis(); - finalCtx.update(entityService.findEntityDataByQuery(tenantId, customerId, finalCtx.getQuery())); + finalCtx.update(); long end = System.currentTimeMillis(); dynamicQueryInvocationCnt.incrementAndGet(); dynamicQueryTimeSpent.addAndGet(end - start); @@ -304,16 +291,16 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc 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()); - ctx.setQuery(cmd.getQuery()); + TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmd.getCmdId()); + ctx.setAndResolveQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; } private TbAlarmDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) { Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); - TbAlarmDataSubCtx ctx = new TbAlarmDataSubCtx(serviceId, wsService, localSubscriptionService, alarmService, sessionRef, cmd.getCmdId()); - ctx.setQuery(cmd.getQuery()); + TbAlarmDataSubCtx ctx = new TbAlarmDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, alarmService, sessionRef, cmd.getCmdId(), maxEntitiesPerAlarmSubscription); + ctx.setAndResolveQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; } @@ -466,7 +453,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private void cleanupAndCancel(TbAbstractDataSubCtx ctx) { if (ctx != null) { ctx.cancelTasks(); - ctx.clearSubscriptions(); + ctx.clearEntitySubscriptions(); } } 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 b42e23bdea..f87220469b 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 @@ -5,7 +5,7 @@ * 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 + * 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, @@ -15,6 +15,9 @@ */ package org.thingsboard.server.service.subscription; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; import lombok.Data; import lombok.Getter; import lombok.Setter; @@ -22,35 +25,55 @@ 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.id.UserId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AbstractDataQuery; +import org.thingsboard.server.common.data.query.ComplexFilterPredicate; +import org.thingsboard.server.common.data.query.DynamicValue; +import org.thingsboard.server.common.data.query.DynamicValueSourceType; 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.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.common.data.query.FilterPredicateType; +import org.thingsboard.server.common.data.query.KeyFilter; +import org.thingsboard.server.common.data.query.KeyFilterPredicate; +import org.thingsboard.server.common.data.query.SimpleKeyFilterPredicate; import org.thingsboard.server.common.data.query.TsValue; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.entity.EntityService; 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.HashSet; import java.util.List; import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ScheduledFuture; @Slf4j @Data -public abstract class TbAbstractDataSubCtx { +public abstract class TbAbstractDataSubCtx> { protected final String serviceId; protected final TelemetryWebSocketService wsService; + protected final EntityService entityService; protected final TbLocalSubscriptionService localSubscriptionService; + protected final AttributesService attributesService; protected final TelemetryWebSocketSessionRef sessionRef; protected final int cmdId; protected final Map subToEntityIdMap; - + protected final Set subToDynamicValueKeySet; + @Getter + protected final Map> dynamicValues; @Getter protected PageData data; @Getter @@ -59,14 +82,209 @@ public abstract class TbAbstractDataSubCtx { @Setter protected volatile ScheduledFuture refreshTask; - public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TbLocalSubscriptionService localSubscriptionService, - TelemetryWebSocketSessionRef sessionRef, int cmdId) { + public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService, + EntityService entityService, TbLocalSubscriptionService localSubscriptionService, + AttributesService attributesService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { this.serviceId = serviceId; this.wsService = wsService; + this.entityService = entityService; this.localSubscriptionService = localSubscriptionService; + this.attributesService = attributesService; this.sessionRef = sessionRef; this.cmdId = cmdId; this.subToEntityIdMap = new HashMap<>(); + this.subToDynamicValueKeySet = new HashSet<>(); + this.dynamicValues = new HashMap<>(); + } + + public void setAndResolveQuery(T query) { + dynamicValues.clear(); + this.query = query; + if (query.getKeyFilters() != null) { + for (KeyFilter filter : query.getKeyFilters()) { + registerDynamicValues(filter.getPredicate()); + } + } + resolve(getTenantId(), getCustomerId(), getUserId()); + } + + public void resolve(TenantId tenantId, CustomerId customerId, UserId userId) { + List> futures = new ArrayList<>(); + for (DynamicValueKey key : dynamicValues.keySet()) { + switch (key.getSourceType()) { + case CURRENT_TENANT: + futures.add(resolveEntityValue(tenantId, tenantId, key)); + break; + case CURRENT_CUSTOMER: + if (customerId != null && !customerId.isNullUid()) { + futures.add(resolveEntityValue(tenantId, customerId, key)); + } + break; + case CURRENT_USER: + if (userId != null && !userId.isNullUid()) { + futures.add(resolveEntityValue(tenantId, userId, key)); + } + break; + } + } + try { + Map> tmpSubMap = new HashMap<>(); + for (DynamicValueKeySub sub : Futures.successfulAsList(futures).get()) { + tmpSubMap.computeIfAbsent(sub.getEntityId(), tmp -> new HashMap<>()).put(sub.getKey().getSourceAttribute(), sub); + } + for (EntityId entityId : tmpSubMap.keySet()) { + Map keyStates = new HashMap<>(); + Map dynamicValueKeySubMap = tmpSubMap.get(entityId); + dynamicValueKeySubMap.forEach((k, v) -> keyStates.put(k, v.getLastUpdateTs())); + int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); + TbAttributeSubscription sub = TbAttributeSubscription.builder() + .serviceId(serviceId) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(subIdx) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityId) + .updateConsumer((s, subscriptionUpdate) -> dynamicValueSubUpdate(s, subscriptionUpdate, dynamicValueKeySubMap)) + .allKeys(false) + .keyStates(keyStates) + .scope(TbAttributeSubscriptionScope.SERVER_SCOPE) + .build(); + subToDynamicValueKeySet.add(subIdx); + localSubscriptionService.addSubscription(sub); + } + + + } catch (InterruptedException | ExecutionException e) { + log.info("[{}][{}][{}] Failed to resolve dynamic values: {}", tenantId, customerId, userId, dynamicValues.keySet()); + } + + } + + private void dynamicValueSubUpdate(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, + Map dynamicValueKeySubMap) { + 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])); + }); + + boolean invalidateFilter = false; + for (Map.Entry entry : latestUpdate.entrySet()) { + String k = entry.getKey(); + TsValue tsValue = entry.getValue(); + DynamicValueKeySub sub = dynamicValueKeySubMap.get(k); + if (sub.updateValue(tsValue)) { + invalidateFilter = true; + updateDynamicValuesByKey(sub, tsValue); + } + } + + if (invalidateFilter) { + update(); + } + } + + public void fetchData() { + this.data = findEntityData(); + } + + protected PageData findEntityData() { + PageData result = entityService.findEntityDataByQuery(getTenantId(), getCustomerId(), buildEntityDataQuery()); + if (log.isTraceEnabled()) { + result.getData().forEach(ed -> { + log.trace("[{}][{}] EntityData: {}", getSessionId(), getCmdId(), ed); + }); + } + return result; + } + + protected abstract void update(); + + protected abstract EntityDataQuery buildEntityDataQuery(); + + public List getEntitiesData() { + return data.getData(); + } + + @Data + private static class DynamicValueKeySub { + private final DynamicValueKey key; + private final EntityId entityId; + private long lastUpdateTs; + private String lastUpdateValue; + + boolean updateValue(TsValue value) { + if (value.getTs() > lastUpdateTs && !lastUpdateValue.equals(value.getValue())) { + this.lastUpdateTs = value.getTs(); + this.lastUpdateValue = value.getValue(); + return true; + } else { + return false; + } + } + } + + private ListenableFuture resolveEntityValue(TenantId tenantId, EntityId entityId, DynamicValueKey key) { + ListenableFuture> entry = attributesService.find(tenantId, entityId, + TbAttributeSubscriptionScope.SERVER_SCOPE.name(), key.getSourceAttribute()); + return Futures.transform(entry, attributeOpt -> { + DynamicValueKeySub sub = new DynamicValueKeySub(key, entityId); + if (attributeOpt.isPresent()) { + AttributeKvEntry attribute = attributeOpt.get(); + sub.setLastUpdateTs(attribute.getLastUpdateTs()); + sub.setLastUpdateValue(attribute.getValueAsString()); + updateDynamicValuesByKey(sub, new TsValue(attribute.getLastUpdateTs(), attribute.getValueAsString())); + } + return sub; + }, MoreExecutors.directExecutor()); + } + + private void updateDynamicValuesByKey(DynamicValueKeySub sub, TsValue tsValue) { + DynamicValueKey dvk = sub.getKey(); + switch (dvk.getPredicateType()) { + case STRING: + dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(tsValue.getValue())); + break; + case NUMERIC: + try { + Double dValue = Double.parseDouble(tsValue.getValue()); + dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(dValue)); + } catch (NumberFormatException e) { + dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(null)); + } + break; + case BOOLEAN: + Boolean bValue = Boolean.parseBoolean(tsValue.getValue()); + dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(bValue)); + break; + } + } + + private void registerDynamicValues(KeyFilterPredicate predicate) { + switch (predicate.getType()) { + case STRING: + case NUMERIC: + case BOOLEAN: + Optional value = getDynamicValueFromSimplePredicate((SimpleKeyFilterPredicate) predicate); + if (value.isPresent()) { + DynamicValue dynamicValue = value.get(); + DynamicValueKey key = new DynamicValueKey( + predicate.getType(), + dynamicValue.getSourceType(), + dynamicValue.getSourceAttribute()); + dynamicValues.computeIfAbsent(key, tmp -> new ArrayList<>()).add(dynamicValue); + } + break; + case COMPLEX: + ((ComplexFilterPredicate) predicate).getPredicates().forEach(this::registerDynamicValues); + } + } + + private Optional> getDynamicValueFromSimplePredicate(SimpleKeyFilterPredicate predicate) { + if (predicate.getValue().getUserValue() == null) { + return Optional.ofNullable(predicate.getValue().getDynamicValue()); + } else { + return Optional.empty(); + } } public String getSessionId() { @@ -81,7 +299,11 @@ public abstract class TbAbstractDataSubCtx { return sessionRef.getSecurityCtx().getCustomerId(); } - public void clearSubscriptions() { + public UserId getUserId() { + return sessionRef.getSecurityCtx().getId(); + } + + public void clearEntitySubscriptions() { if (subToEntityIdMap != null) { for (Integer subId : subToEntityIdMap.keySet()) { localSubscriptionService.cancelSubscription(sessionRef.getSessionId(), subId); @@ -90,10 +312,6 @@ public abstract class TbAbstractDataSubCtx { } } - public void setData(PageData data) { - this.data = data; - } - public void setRefreshTask(ScheduledFuture task) { this.refreshTask = task; } @@ -204,4 +422,15 @@ public abstract class TbAbstractDataSubCtx { } abstract void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues); + + @Data + private static class DynamicValueKey { + @Getter + private final FilterPredicateType predicateType; + @Getter + private final DynamicValueSourceType sourceType; + @Getter + private final String sourceAttribute; + } + } 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 db849ef5d5..66fad14770 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,17 +27,22 @@ 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.EntityDataPageLink; +import org.thingsboard.server.common.data.query.EntityDataQuery; +import org.thingsboard.server.common.data.query.EntityDataSortOrder; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.alarm.AlarmService; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.entity.EntityService; +import org.thingsboard.server.dao.model.ModelConstants; 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; @@ -58,7 +63,8 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { @Setter private final HashMap alarmsMap; - private final List alarmSubscriptions; + private final int maxEntitiesPerAlarmSubscription; + @Getter @Setter private PageData alarms; @@ -67,14 +73,14 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { private boolean tooManyEntities; public TbAlarmDataSubCtx(String serviceId, TelemetryWebSocketService wsService, - TbLocalSubscriptionService localSubscriptionService, - AlarmService alarmService, - TelemetryWebSocketSessionRef sessionRef, int cmdId) { - super(serviceId, wsService, localSubscriptionService, sessionRef, cmdId); + EntityService entityService, TbLocalSubscriptionService localSubscriptionService, + AttributesService attributesService, AlarmService alarmService, + TelemetryWebSocketSessionRef sessionRef, int cmdId, int maxEntitiesPerAlarmSubscription) { + super(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmdId); + this.maxEntitiesPerAlarmSubscription = maxEntitiesPerAlarmSubscription; this.alarmService = alarmService; this.entitiesMap = new LinkedHashMap<>(); this.alarmsMap = new HashMap<>(); - this.alarmSubscriptions = new ArrayList<>(); } public void fetchAlarms() { @@ -85,8 +91,8 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { wsService.sendWsMsg(getSessionId(), update); } - public void setData(PageData data) { - super.setData(data); + public void fetchData() { + super.fetchData(); entitiesMap.clear(); tooManyEntities = data.hasNext(); for (EntityData entityData : data.getData()) { @@ -233,4 +239,22 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { } return true; } + + @Override + protected void update() { + + } + + @Override + protected EntityDataQuery buildEntityDataQuery() { + EntityDataSortOrder sortOrder = query.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(maxEntitiesPerAlarmSubscription, 0, null, entitiesSortOrder); + return new EntityDataQuery(query.getEntityFilter(), edpl, query.getEntityFields(), query.getLatestValues(), query.getKeyFilters()); + } } 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 f57de8e504..e600337cb9 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 @@ -27,6 +27,8 @@ import org.thingsboard.server.common.data.query.EntityDataQuery; 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.attributes.AttributesService; +import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; @@ -62,10 +64,10 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { private TimeSeriesCmd curTsCmd; private LatestValueCmd latestValueCmd; - public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, - TbLocalSubscriptionService localSubscriptionService, + public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService, + TbLocalSubscriptionService localSubscriptionService, AttributesService attributesService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { - super(serviceId, wsService, localSubscriptionService, sessionRef, cmdId); + super(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmdId); } @Override @@ -169,7 +171,8 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null); } - public void update(PageData newData) { + public synchronized void update() { + PageData newData = findEntityData(); Map oldDataMap; if (data != null && !data.getData().isEmpty()) { oldDataMap = data.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); @@ -226,4 +229,9 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { curTsCmd = cmd.getTsCmd(); latestValueCmd = cmd.getLatestCmd(); } + + @Override + protected EntityDataQuery buildEntityDataQuery() { + return query; + } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java index 36a5aa37ad..325b1ba11f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/AbstractDataQuery.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.common.data.query; -import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Getter; import lombok.ToString; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/BooleanFilterPredicate.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/BooleanFilterPredicate.java index 63d7cb2d9a..7fdea11d83 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/BooleanFilterPredicate.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/BooleanFilterPredicate.java @@ -18,7 +18,7 @@ package org.thingsboard.server.common.data.query; import lombok.Data; @Data -public class BooleanFilterPredicate implements KeyFilterPredicate { +public class BooleanFilterPredicate implements SimpleKeyFilterPredicate { private BooleanOperation operation; private FilterPredicateValue value; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/NumericFilterPredicate.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/NumericFilterPredicate.java index 3d0789e80e..22ca8b85a4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/NumericFilterPredicate.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/NumericFilterPredicate.java @@ -18,7 +18,7 @@ package org.thingsboard.server.common.data.query; import lombok.Data; @Data -public class NumericFilterPredicate implements KeyFilterPredicate { +public class NumericFilterPredicate implements SimpleKeyFilterPredicate { private NumericOperation operation; private FilterPredicateValue value; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/SimpleKeyFilterPredicate.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/SimpleKeyFilterPredicate.java new file mode 100644 index 0000000000..0d796b90d5 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/SimpleKeyFilterPredicate.java @@ -0,0 +1,22 @@ +/** + * 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; + +public interface SimpleKeyFilterPredicate extends KeyFilterPredicate { + + FilterPredicateValue getValue(); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/StringFilterPredicate.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/StringFilterPredicate.java index 36a47b0675..77459529d0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/StringFilterPredicate.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/StringFilterPredicate.java @@ -18,7 +18,7 @@ package org.thingsboard.server.common.data.query; import lombok.Data; @Data -public class StringFilterPredicate implements KeyFilterPredicate { +public class StringFilterPredicate implements SimpleKeyFilterPredicate { private StringOperation operation; private FilterPredicateValue value;