Browse Source

Alarm Implementation

pull/3068/head
Andrii Shvaika 6 years ago
parent
commit
d102ef7d3f
  1. 30
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 33
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java
  3. 57
      application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java
  4. 31
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java
  5. 3
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubscriptionService.java
  6. 9
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java
  7. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java
  8. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/DataUpdate.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java
  11. 24
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/UnsubscribeCmd.java
  12. 4
      application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java
  13. 20
      dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
  14. 57
      dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractSearchTimeDao.java
  15. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java
  16. 42
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java
  17. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  18. 8
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java
  19. 2
      dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java

30
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<EntityData> entitiesData = entityService.findEntityDataByQuery(ctx.getTenantId(), ctx.getCustomerId(), edq);
List<EntityData> 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<Integer> 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<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>());
TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, sessionRef, cmd.getCmdId());
TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, 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);
}
}

33
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<T extends AbstractDataQuery> {
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<Integer, EntityId> 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<T extends AbstractDataQuery> {
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);
}
}
}

57
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<AlarmDataQuery> {
private final TbLocalSubscriptionService localSubscriptionService;
private final AlarmService alarmService;
@Getter
@Setter
@ -62,31 +57,22 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
@Setter
private boolean tooManyEntities;
private Map<Integer, EntityId> 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<AlarmData> 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<TbSubscription> subscriptions = createSubscriptions();
subscriptions.forEach(localSubscriptionService::addSubscription);
}
}
public void setEntitiesData(PageData<EntityData> entitiesData) {
@ -117,16 +103,16 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
return this.alarms;
}
public List<TbSubscription> createSubscriptions() {
public void createSubscriptions() {
clearSubscriptions();
this.subToEntityIdMap = new HashMap<>();
AlarmDataPageLink pageLink = query.getPageLink();
long startTs = System.currentTimeMillis() - pageLink.getTimeWindow();
List<TbSubscription> 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<AlarmDataQuery> {
.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<AlarmDataQuery> {
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<AlarmDataQuery> {
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();

31
application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java

@ -61,13 +61,13 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
@Getter
@Setter
private boolean initialDataSent;
private Map<Integer, EntityId> 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<EntityData> data) {
@ -272,28 +272,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
private EntityData getDataForEntity(EntityId entityId) {
return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null);
}
public Collection<Integer> clearSubscriptions() {
if (subToEntityIdMap != null) {
List<Integer> 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<EntityData> newData) {
Map<EntityId, EntityData> oldDataMap;
if (data != null && !data.getData().isEmpty()) {

3
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);

9
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);

2
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;

2
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<T> {
private final int cmdId;

2
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;

2
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<EntityData> {
public EntityDataUpdate(int cmdId, PageData<EntityData> data, List<EntityData> update) {

24
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();
}

4
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 {

20
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 <E extends HasName, I extends EntityId> ListenableFuture<List<Void>>
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 <E extends HasName, I extends EntityId> 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<String> 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);

57
dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractSearchTimeDao.java

@ -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<E extends BaseEntity<D>, D> extends JpaAbstractDao<E, D> {
//TODO 3.1: refactoring to createdTime column
public static <T> Specification<T> getTimeSearchPageSpec(TimePageLink pageLink, String idColumn) {
return new Specification<T>() {
@Override
public Predicate toPredicate(Root<T> root, CriteriaQuery<?> criteriaQuery, CriteriaBuilder criteriaBuilder) {
List<Predicate> 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]));
}
};
}
}

5
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<AlarmEntity, UUID> {
"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<AlarmEntity, UUID> {
"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<AlarmEntity, UUID> {
@Param("relationType") String relationType,
@Param("startTime") Long startTime,
@Param("endTime") Long endTime,
@Param("alarmStatuses") Set<AlarmStatus> alarmStatuses,
@Param("searchText") String searchText,
Pageable pageable);

42
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<AlarmEntity, Alarm> implements A
public PageData<AlarmInfo> 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<AlarmStatus> 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())
)
);
}

4
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<EventEntity, Event> implements EventDao {
public class JpaBaseEventDao extends JpaAbstractDao<EventEntity, Event> implements EventDao {
private final UUID systemTenantId = NULL_UUID;

8
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;

2
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 {

Loading…
Cancel
Save