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 de4fef7928..5b83fc5212 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -34,7 +34,6 @@ import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataPageLink; @@ -55,19 +54,18 @@ import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; -import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityHistoryCmd; import org.thingsboard.server.service.telemetry.cmd.v2.GetTsCmd; import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; import org.thingsboard.server.service.telemetry.cmd.v2.TimeSeriesCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.UnsubscribeCmd; import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.Arrays; -import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.LinkedHashMap; @@ -203,7 +201,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }); } ctx.setData(data); - ctx.cancelRefreshTask(); + ctx.cancelTasks(); 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; @@ -262,11 +260,20 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc PageData entitiesData = entityService.findEntityDataByQuery(ctx.getTenantId(), ctx.getCustomerId(), edq); List entities = entitiesData.getData(); ctx.setEntitiesData(entitiesData); + ctx.cancelTasks(); if (entities.isEmpty()) { AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(Collections.emptyList(), 1, 0, false), null); wsService.sendWsMsg(ctx.getSessionId(), update); } else { - ctx.fetchAlarmsAndCreateSubscriptions(); + ctx.fetchAlarms(); + if (adq.getPageLink().getTimeWindow() > 0) { + ctx.createSubscriptions(); + TbAlarmDataSubCtx finalCtx = ctx; + ScheduledFuture task = scheduler.scheduleWithFixedDelay( + finalCtx::cleanupOldAlarms, + dynamicPageLinkRefreshInterval, dynamicPageLinkRefreshInterval, TimeUnit.SECONDS); + finalCtx.setRefreshTask(task); + } } } @@ -297,14 +304,13 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } - private void clearSubs(TbEntityDataSubCtx ctx) { - Collection oldSubIds = ctx.clearSubscriptions(); - oldSubIds.forEach(subId -> localSubscriptionService.cancelSubscription(serviceId, subId)); + 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, sessionRef, cmd.getCmdId()); + TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, localSubscriptionService, sessionRef, cmd.getCmdId()); ctx.setQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; @@ -468,13 +474,13 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } @Override - public void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd cmd) { + public void cancelSubscription(String sessionId, UnsubscribeCmd cmd) { cleanupAndCancel(getSubCtx(sessionId, cmd.getCmdId())); } - private void cleanupAndCancel(TbEntityDataSubCtx ctx) { + private void cleanupAndCancel(TbAbstractDataSubCtx ctx) { if (ctx != null) { - ctx.cancelRefreshTask(); + ctx.cancelTasks(); clearSubs(ctx); } } 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 b323ab7240..244613a984 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 @@ -20,26 +20,37 @@ 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.EntityId; 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; +import java.util.Collection; +import java.util.Map; +import java.util.concurrent.ScheduledFuture; + @Slf4j @Data public abstract class TbAbstractDataSubCtx { protected final String serviceId; protected final TelemetryWebSocketService wsService; + protected final TbLocalSubscriptionService localSubscriptionService; protected final TelemetryWebSocketSessionRef sessionRef; protected final int cmdId; @Getter @Setter protected T query; + protected Map subToEntityIdMap; + @Setter + protected volatile ScheduledFuture refreshTask; - public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { + public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TbLocalSubscriptionService localSubscriptionService, + TelemetryWebSocketSessionRef sessionRef, int cmdId) { this.serviceId = serviceId; this.wsService = wsService; + this.localSubscriptionService = localSubscriptionService; this.sessionRef = sessionRef; this.cmdId = cmdId; } @@ -56,4 +67,24 @@ public abstract class TbAbstractDataSubCtx { return sessionRef.getSecurityCtx().getCustomerId(); } + public void clearSubscriptions(){ + if (subToEntityIdMap != null) { + for (Integer subId : subToEntityIdMap.keySet()) { + localSubscriptionService.cancelSubscription(sessionRef.getSessionId(), subId); + } + subToEntityIdMap.clear(); + } + } + + public void setRefreshTask(ScheduledFuture task) { + this.refreshTask = task; + } + + public void cancelTasks() { + if (this.refreshTask != null) { + log.trace("[{}][{}] Canceling old refresh task", sessionRef.getSessionId(), cmdId); + this.refreshTask.cancel(true); + } + } + } 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 c8f4d0f2e6..c00450e871 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 @@ -30,24 +30,19 @@ import org.thingsboard.server.common.data.query.EntityData; 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.AlarmDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate; -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; @Slf4j public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { - private final TbLocalSubscriptionService localSubscriptionService; private final AlarmService alarmService; @Getter @Setter @@ -62,31 +57,22 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { @Setter private boolean tooManyEntities; - private Map subToEntityIdMap; - public TbAlarmDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TbLocalSubscriptionService localSubscriptionService, AlarmService alarmService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { - super(serviceId, wsService, sessionRef, cmdId); - this.localSubscriptionService = localSubscriptionService; + super(serviceId, wsService, localSubscriptionService, sessionRef, cmdId); this.alarmService = alarmService; this.entitiesMap = new LinkedHashMap<>(); this.alarmsMap = new HashMap<>(); } - public void fetchAlarmsAndCreateSubscriptions() { + public void fetchAlarms() { PageData alarms = alarmService.findAlarmDataByQueryForEntities(getTenantId(), getCustomerId(), query.getPageLink(), getOrderedEntityIds()); alarms = setAndMergeAlarmsData(alarms); AlarmDataUpdate update = new AlarmDataUpdate(cmdId, alarms, null); wsService.sendWsMsg(getSessionId(), update); - if (query.getPageLink().getTimeWindow() > 0) { - clearSubscriptions(); - //TODO: refresh list of entities periodically (similar to time-series subscription). - List subscriptions = createSubscriptions(); - subscriptions.forEach(localSubscriptionService::addSubscription); - } } public void setEntitiesData(PageData entitiesData) { @@ -117,16 +103,16 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { return this.alarms; } - public List createSubscriptions() { + public void createSubscriptions() { + clearSubscriptions(); this.subToEntityIdMap = new HashMap<>(); AlarmDataPageLink pageLink = query.getPageLink(); long startTs = System.currentTimeMillis() - pageLink.getTimeWindow(); - List result = new ArrayList<>(); for (EntityData entityData : entitiesMap.values()) { int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); subToEntityIdMap.put(subIdx, entityData.getEntityId()); log.trace("[{}][{}][{}] Creating alarms subscription for [{}] with query: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), pageLink); - result.add(TbAlarmsSubscription.builder() + TbAlarmsSubscription subscription = TbAlarmsSubscription.builder() .type(TbSubscriptionType.ALARMS) .serviceId(serviceId) .sessionId(sessionRef.getSessionId()) @@ -135,17 +121,8 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { .entityId(entityData.getEntityId()) .updateConsumer(this::sendWsMsg) .ts(startTs) - .build()); - } - return result; - } - - public void clearSubscriptions() { - if (subToEntityIdMap != null) { - for (Integer subId : subToEntityIdMap.keySet()) { - localSubscriptionService.cancelSubscription(sessionRef.getSessionId(), subId); - } - subToEntityIdMap.clear(); + .build(); + localSubscriptionService.addSubscription(subscription); } } @@ -155,7 +132,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { if (subscriptionUpdate.isAlarmDeleted()) { Alarm deleted = alarmsMap.remove(alarmId); if (deleted != null) { - fetchAlarmsAndCreateSubscriptions(); + fetchAlarms(); } } else { AlarmData current = alarmsMap.get(alarmId); @@ -167,14 +144,28 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { alarmsMap.put(alarmId, updated); wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated))); } else { - fetchAlarmsAndCreateSubscriptions(); + fetchAlarms(); } } else if (matchesFilter && query.getPageLink().getPage() == 0) { - fetchAlarmsAndCreateSubscriptions(); + fetchAlarms(); } } } + public void cleanupOldAlarms() { + long expTime = System.currentTimeMillis() - query.getPageLink().getTimeWindow(); + boolean shouldRefresh = false; + for (AlarmData alarmData : alarms.getData()) { + if (alarmData.getCreatedTime() < expTime) { + shouldRefresh = true; + break; + } + } + if (shouldRefresh) { + fetchAlarms(); + } + } + private boolean filter(Alarm alarm) { AlarmDataPageLink filter = query.getPageLink(); long startTs = System.currentTimeMillis() - filter.getTimeWindow(); 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 85318efe26..80505b685b 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 @@ -61,13 +61,13 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { @Getter @Setter private boolean initialDataSent; - private Map subToEntityIdMap; - private volatile ScheduledFuture refreshTask; private TimeSeriesCmd curTsCmd; private LatestValueCmd latestValueCmd; - public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { - super(serviceId, wsService, sessionRef, cmdId); + public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, + TbLocalSubscriptionService localSubscriptionService, + TelemetryWebSocketSessionRef sessionRef, int cmdId) { + super(serviceId, wsService, localSubscriptionService, sessionRef, cmdId); } public void setData(PageData data) { @@ -272,28 +272,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { private EntityData getDataForEntity(EntityId entityId) { return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null); } - - public Collection clearSubscriptions() { - if (subToEntityIdMap != null) { - List oldSubIds = new ArrayList<>(subToEntityIdMap.keySet()); - subToEntityIdMap.clear(); - return oldSubIds; - } else { - return Collections.emptyList(); - } - } - - public void setRefreshTask(ScheduledFuture task) { - this.refreshTask = task; - } - - public void cancelRefreshTask() { - if (this.refreshTask != null) { - log.trace("[{}][{}] Canceling old refresh task", sessionRef.getSessionId(), cmdId); - this.refreshTask.cancel(true); - } - } - + public TbEntityDataSubCtxUpdateResult update(PageData newData) { Map oldDataMap; if (data != null && !data.getData().isEmpty()) { diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java index b805cf576d..4f5d9543b9 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java @@ -19,6 +19,7 @@ import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.UnsubscribeCmd; public interface TbEntityDataSubscriptionService { @@ -26,7 +27,7 @@ public interface TbEntityDataSubscriptionService { void handleCmd(TelemetryWebSocketSessionRef sessionId, AlarmDataCmd cmd); - void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd subscriptionId); + void cancelSubscription(String sessionId, UnsubscribeCmd subscriptionId); void cancelAllSessionSubscriptions(String sessionId); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java index 7cdb911ae2..00540ce02b 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java @@ -63,10 +63,12 @@ import org.thingsboard.server.service.telemetry.cmd.v1.TelemetryPluginCmd; import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.telemetry.cmd.v1.TimeseriesSubscriptionCmd; import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; +import org.thingsboard.server.service.telemetry.cmd.v2.UnsubscribeCmd; import org.thingsboard.server.service.telemetry.exception.UnauthorizedException; import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; @@ -214,7 +216,10 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi cmdsWrapper.getAlarmDataCmds().forEach(cmd -> handleWsAlarmDataCmd(sessionRef, cmd)); } if (cmdsWrapper.getEntityDataUnsubscribeCmds() != null) { - cmdsWrapper.getEntityDataUnsubscribeCmds().forEach(cmd -> handleWsEntityDataUnsubscribeCmd(sessionRef, cmd)); + cmdsWrapper.getEntityDataUnsubscribeCmds().forEach(cmd -> handleWsDataUnsubscribeCmd(sessionRef, cmd)); + } + if (cmdsWrapper.getAlarmDataUnsubscribeCmds() != null) { + cmdsWrapper.getAlarmDataUnsubscribeCmds().forEach(cmd -> handleWsDataUnsubscribeCmd(sessionRef, cmd)); } } } catch (IOException e) { @@ -244,7 +249,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } } - private void handleWsEntityDataUnsubscribeCmd(TelemetryWebSocketSessionRef sessionRef, EntityDataUnsubscribeCmd cmd) { + private void handleWsDataUnsubscribeCmd(TelemetryWebSocketSessionRef sessionRef, UnsubscribeCmd cmd) { String sessionId = sessionRef.getSessionId(); log.debug("[{}] Processing: {}", sessionId, cmd); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java index 6782c2f595..b886cff349 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java @@ -18,7 +18,7 @@ package org.thingsboard.server.service.telemetry.cmd.v2; import lombok.Data; @Data -public class AlarmDataUnsubscribeCmd { +public class AlarmDataUnsubscribeCmd implements UnsubscribeCmd { private final int cmdId; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java index 81dc4ef902..c841e220fb 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.telemetry.cmd.v2; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import lombok.AllArgsConstructor; import lombok.Data; import org.thingsboard.server.common.data.page.PageData; @@ -24,6 +25,7 @@ import java.util.List; @Data @AllArgsConstructor +@JsonIgnoreProperties(ignoreUnknown = true) public abstract class DataUpdate { private final int cmdId; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java index 4af4106a9c..f75f622db1 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java @@ -18,7 +18,7 @@ package org.thingsboard.server.service.telemetry.cmd.v2; import lombok.Data; @Data -public class EntityDataUnsubscribeCmd { +public class EntityDataUnsubscribeCmd implements UnsubscribeCmd { private final int cmdId; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java index 28cf6c8719..81a1ee1778 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.telemetry.cmd.v2; import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonProperty; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; @@ -23,6 +24,7 @@ import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; import java.util.List; + public class EntityDataUpdate extends DataUpdate { public EntityDataUpdate(int cmdId, PageData data, List update) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/UnsubscribeCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/UnsubscribeCmd.java new file mode 100644 index 0000000000..a2f5bfe770 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/UnsubscribeCmd.java @@ -0,0 +1,24 @@ +/** + * 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; + +public interface UnsubscribeCmd { + + int getCmdId(); + +} diff --git a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java index 15da972cf5..0a5dae47b7 100644 --- a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java @@ -26,9 +26,9 @@ import java.util.Arrays; @RunWith(ClasspathSuite.class) @ClasspathSuite.ClassnameFilters({ -// "org.thingsboard.server.controller.sql.WebsocketApiSqlTest", + "org.thingsboard.server.controller.sql.WebsocketApiSqlTest", // "org.thingsboard.server.controller.sql.EntityQueryControllerSqlTest", - "org.thingsboard.server.controller.sql.*Test", +// "org.thingsboard.server.controller.sql.*Test", }) public class ControllerSqlTestSuite { diff --git a/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java index c308ebe322..0392ec7f25 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java @@ -52,6 +52,7 @@ import org.thingsboard.server.dao.service.DataValidator; import java.io.PrintWriter; import java.io.StringWriter; import java.util.List; +import java.util.UUID; import static org.thingsboard.server.dao.service.Validator.validateEntityId; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -111,8 +112,8 @@ public class AuditLogServiceImpl implements AuditLogService { @Override public ListenableFuture> - logEntityAction(TenantId tenantId, CustomerId customerId, UserId userId, String userName, I entityId, E entity, - ActionType actionType, Exception e, Object... additionalInfo) { + logEntityAction(TenantId tenantId, CustomerId customerId, UserId userId, String userName, I entityId, E entity, + ActionType actionType, Exception e, Object... additionalInfo) { if (canLog(entityId.getEntityType(), actionType)) { JsonNode actionData = constructActionData(entityId, entity, actionType, additionalInfo); ActionStatus actionStatus = ActionStatus.SUCCESS; @@ -123,7 +124,8 @@ public class AuditLogServiceImpl implements AuditLogService { } else { try { entityName = entityService.fetchEntityNameAsync(tenantId, entityId).get(); - } catch (Exception ex) {} + } catch (Exception ex) { + } } if (e != null) { actionStatus = ActionStatus.FAILURE; @@ -152,10 +154,10 @@ public class AuditLogServiceImpl implements AuditLogService { } private JsonNode constructActionData(I entityId, E entity, - ActionType actionType, - Object... additionalInfo) { + ActionType actionType, + Object... additionalInfo) { ObjectNode actionData = objectMapper.createObjectNode(); - switch(actionType) { + switch (actionType) { case ADDED: case UPDATED: case ALARM_ACK: @@ -202,7 +204,7 @@ public class AuditLogServiceImpl implements AuditLogService { scope = extractParameter(String.class, 0, additionalInfo); actionData.put("scope", scope); List keys = extractParameter(List.class, 1, additionalInfo); - ArrayNode attrsArrayNode = actionData.putArray("attributes"); + ArrayNode attrsArrayNode = actionData.putArray("attributes"); if (keys != null) { keys.forEach(attrsArrayNode::add); } @@ -294,7 +296,9 @@ public class AuditLogServiceImpl implements AuditLogService { ActionStatus actionStatus, String actionFailureDetails) { AuditLog result = new AuditLog(); - result.setId(new AuditLogId(Uuids.timeBased())); + UUID id = Uuids.timeBased(); + result.setId(new AuditLogId(id)); + result.setCreatedTime(id.timestamp()); result.setTenantId(tenantId); result.setEntityId(entityId); result.setEntityName(entityName); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractSearchTimeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractSearchTimeDao.java deleted file mode 100644 index 485a069674..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractSearchTimeDao.java +++ /dev/null @@ -1,57 +0,0 @@ -/** - * 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; - -import com.datastax.oss.driver.api.core.uuid.Uuids; -import org.springframework.data.jpa.domain.Specification; -import org.thingsboard.server.common.data.UUIDConverter; -import org.thingsboard.server.common.data.page.TimePageLink; -import org.thingsboard.server.dao.model.BaseEntity; - -import javax.persistence.criteria.CriteriaBuilder; -import javax.persistence.criteria.CriteriaQuery; -import javax.persistence.criteria.Predicate; -import javax.persistence.criteria.Root; -import java.util.ArrayList; -import java.util.List; -import java.util.UUID; - -/** - * Created by Valerii Sosliuk on 5/4/2017. - */ -public abstract class JpaAbstractSearchTimeDao, D> extends JpaAbstractDao { - - //TODO 3.1: refactoring to createdTime column - public static Specification getTimeSearchPageSpec(TimePageLink pageLink, String idColumn) { - return new Specification() { - @Override - public Predicate toPredicate(Root root, CriteriaQuery criteriaQuery, CriteriaBuilder criteriaBuilder) { - List predicates = new ArrayList<>(); - if (pageLink.getStartTime() != null) { - UUID startOf = Uuids.startOf(pageLink.getStartTime()); - Predicate lowerBound = criteriaBuilder.greaterThanOrEqualTo(root.get(idColumn), UUIDConverter.fromTimeUUID(startOf)); - predicates.add(lowerBound); - } - if (pageLink.getEndTime() != null) { - UUID endOf = Uuids.endOf(pageLink.getEndTime()); - Predicate upperBound = criteriaBuilder.lessThanOrEqualTo(root.get(idColumn), UUIDConverter.fromTimeUUID(endOf)); - predicates.add(upperBound); - } - return criteriaBuilder.and(predicates.toArray(new Predicate[0])); - } - }; - } -} 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 8fd5113ed1..883da98050 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java @@ -20,6 +20,7 @@ import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.query.Param; +import org.thingsboard.server.common.data.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; @@ -32,6 +33,7 @@ import org.thingsboard.server.dao.util.SqlDao; import java.util.Collection; import java.util.List; +import java.util.Set; import java.util.UUID; /** @@ -55,6 +57,7 @@ public interface AlarmRepository extends CrudRepository { "AND re.fromType = :affectedEntityType " + "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 (LOWER(a.type) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.severity) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.status) LIKE LOWER(CONCAT(:searchText, '%')))", @@ -68,6 +71,7 @@ public interface AlarmRepository extends CrudRepository { "AND re.fromType = :affectedEntityType " + "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 (LOWER(a.type) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.severity) LIKE LOWER(CONCAT(:searchText, '%'))" + "OR LOWER(a.status) LIKE LOWER(CONCAT(:searchText, '%')))") @@ -77,6 +81,7 @@ public interface AlarmRepository extends CrudRepository { @Param("relationType") String relationType, @Param("startTime") Long startTime, @Param("endTime") Long endTime, + @Param("alarmStatuses") Set 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 9df831c8f6..463e439174 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 @@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmQuery; import org.thingsboard.server.common.data.alarm.AlarmSearchStatus; +import org.thingsboard.server.common.data.alarm.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; @@ -40,8 +41,10 @@ import org.thingsboard.server.dao.sql.query.AlarmQueryRepository; import org.thingsboard.server.dao.util.SqlDao; import java.util.Collection; +import java.util.Collections; import java.util.List; import java.util.Objects; +import java.util.Set; import java.util.UUID; /** @@ -96,29 +99,24 @@ public class JpaAlarmDao extends JpaAbstractDao implements A public PageData findAlarms(TenantId tenantId, AlarmQuery query) { log.trace("Try to find alarms by entity [{}], status [{}] and pageLink [{}]", query.getAffectedEntityId(), query.getStatus(), query.getPageLink()); EntityId affectedEntity = query.getAffectedEntityId(); - - //TODO 3.1: add search by statuses -// String searchStatusName; -// if (query.getSearchStatus() == null && query.getStatus() == null) { -// searchStatusName = AlarmSearchStatus.ANY.name(); -// } else if (query.getSearchStatus() != null) { -// searchStatusName = query.getSearchStatus().name(); -// } else { -// searchStatusName = query.getStatus().name(); -// } -// String relationType = BaseAlarmService.ALARM_RELATION_PREFIX; - + Set statusSet = null; + if (query.getSearchStatus() != null) { + statusSet = query.getSearchStatus().getStatuses(); + } else if (query.getStatus() != null){ + statusSet = Collections.singleton(query.getStatus()); + } return DaoUtil.toPageData( - alarmRepository.findAlarms( - tenantId.getId(), - affectedEntity.getId(), - affectedEntity.getEntityType().name(), - AlarmSearchStatus.ANY.name(), - query.getPageLink().getStartTime(), - query.getPageLink().getEndTime(), - Objects.toString(query.getPageLink().getTextSearch(), ""), - DaoUtil.toPageable(query.getPageLink()) - ) + alarmRepository.findAlarms( + tenantId.getId(), + affectedEntity.getId(), + affectedEntity.getEntityType().name(), + AlarmSearchStatus.ANY.name(), + query.getPageLink().getStartTime(), + query.getPageLink().getEndTime(), + statusSet, + Objects.toString(query.getPageLink().getTextSearch(), ""), + DaoUtil.toPageable(query.getPageLink()) + ) ); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java index 18537f90d9..27f9738e16 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java @@ -33,7 +33,7 @@ import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.event.EventDao; import org.thingsboard.server.dao.model.sql.EventEntity; -import org.thingsboard.server.dao.sql.JpaAbstractSearchTimeDao; +import org.thingsboard.server.dao.sql.JpaAbstractDao; import org.thingsboard.server.dao.util.SqlDao; import javax.persistence.criteria.Predicate; @@ -51,7 +51,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; @Slf4j @Component @SqlDao -public class JpaBaseEventDao extends JpaAbstractSearchTimeDao implements EventDao { +public class JpaBaseEventDao extends JpaAbstractDao implements EventDao { private final UUID systemTenantId = NULL_UUID; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java index 306f2c91a9..b96d3db110 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java @@ -18,18 +18,11 @@ package org.thingsboard.server.dao.sql.relation; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.data.domain.PageRequest; -import org.springframework.data.domain.Pageable; -import org.springframework.data.domain.Sort; import org.springframework.data.jpa.domain.Specification; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.UUIDConverter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.SortOrder; -import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.dao.DaoUtil; @@ -37,7 +30,6 @@ import org.thingsboard.server.dao.model.sql.RelationCompositeKey; import org.thingsboard.server.dao.model.sql.RelationEntity; import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; -import org.thingsboard.server.dao.sql.JpaAbstractSearchTimeDao; import org.thingsboard.server.dao.util.SqlDao; import javax.persistence.criteria.Predicate; diff --git a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java index db87deab93..a6ef3935b0 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java @@ -24,7 +24,7 @@ import java.util.Arrays; @RunWith(ClasspathSuite.class) @ClassnameFilters({ - "org.thingsboard.server.dao.service.sql.AlarmServiceSqlTest" + "org.thingsboard.server.dao.service.sql.*SqlTest" }) public class SqlDaoServiceTestSuite {