From 9e1d86d7e8c4837738d5ce696c5cc88b78edee16 Mon Sep 17 00:00:00 2001 From: Viacheslav Klimov Date: Mon, 31 May 2021 15:41:16 +0300 Subject: [PATCH] Refactor --- .../server/controller/AlarmController.java | 5 +- .../server/controller/BaseController.java | 212 +-------------- .../action/RuleEngineEntityActionService.java | 256 ++++++++++++++++++ .../service/ttl/AbstractCleanUpService.java | 6 +- .../ttl/alarms/AlarmsCleanUpService.java | 67 +++-- .../ttl/edge/EdgeEventsCleanUpService.java | 2 +- .../ttl/events/EventsCleanUpService.java | 2 +- .../AbstractTimeseriesCleanUpService.java | 2 +- .../src/main/resources/thingsboard.yml | 2 +- .../server/common/data/DataConstants.java | 1 + .../server/common/data/TenantProfile.java | 1 + .../server/common/data/audit/ActionType.java | 1 + .../server/dao/alarm/AlarmDao.java | 3 +- .../server/dao/sql/alarm/AlarmRepository.java | 5 +- .../server/dao/sql/alarm/JpaAlarmDao.java | 6 +- .../rule/engine/profile/DeviceState.java | 8 + 16 files changed, 341 insertions(+), 238 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java diff --git a/application/src/main/java/org/thingsboard/server/controller/AlarmController.java b/application/src/main/java/org/thingsboard/server/controller/AlarmController.java index 03883fc167..f0af59c7bc 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AlarmController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AlarmController.java @@ -110,8 +110,11 @@ public class AlarmController extends BaseController { checkParameter(ALARM_ID, strAlarmId); try { AlarmId alarmId = new AlarmId(toUUID(strAlarmId)); - checkAlarmId(alarmId, Operation.WRITE); + Alarm alarm = checkAlarmId(alarmId, Operation.WRITE); + logEntityAction(alarm.getOriginator(), alarm, + getCurrentUser().getCustomerId(), + ActionType.ALARM_DELETE, null); sendEntityNotificationMsg(getTenantId(), alarmId, EdgeEventActionType.DELETED); return alarmService.deleteAlarm(getTenantId(), alarmId); diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 505416f4d6..5ccb763445 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -17,7 +17,6 @@ package org.thingsboard.server.controller; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.Getter; import lombok.extern.slf4j.Slf4j; @@ -31,7 +30,6 @@ import org.springframework.web.bind.annotation.ExceptionHandler; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.DashboardInfo; -import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceInfo; import org.thingsboard.server.common.data.DeviceProfile; @@ -79,10 +77,6 @@ import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetsBundleId; -import org.thingsboard.server.common.data.kv.AttributeKvEntry; -import org.thingsboard.server.common.data.kv.DataType; -import org.thingsboard.server.common.data.kv.KvEntry; -import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; @@ -94,9 +88,6 @@ import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.data.widget.WidgetTypeDetails; import org.thingsboard.server.common.data.widget.WidgetsBundle; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.TbMsgDataType; -import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.audit.AuditLogService; @@ -127,6 +118,7 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.action.RuleEngineEntityActionService; import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.firmware.FirmwareStateService; import org.thingsboard.server.service.edge.EdgeNotificationService; @@ -147,11 +139,9 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.mail.MessagingException; import javax.servlet.http.HttpServletResponse; import java.util.List; -import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.UUID; -import java.util.stream.Collectors; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -281,6 +271,9 @@ public abstract class BaseController { @Autowired(required = false) protected EdgeGrpcService edgeGrpcService; + @Autowired + protected RuleEngineEntityActionService ruleEngineEntityActionService; + @Value("${server.log_controller_error_stack_trace}") @Getter private boolean logControllerErrorStackTrace; @@ -811,7 +804,7 @@ public abstract class BaseController { customerId = user.getCustomerId(); } if (e == null) { - pushEntityActionToRuleEngine(entityId, entity, user, customerId, actionType, additionalInfo); + ruleEngineEntityActionService.pushEntityActionToRuleEngine(entityId, entity, user.getTenantId(), customerId, actionType, user, additionalInfo); } auditLogService.logEntityAction(user.getTenantId(), customerId, user.getId(), user.getName(), entityId, entity, actionType, e, additionalInfo); } @@ -821,184 +814,6 @@ public abstract class BaseController { return error != null ? (Exception.class.isInstance(error) ? (Exception) error : new Exception(error)) : null; } - private void pushEntityActionToRuleEngine(I entityId, E entity, User user, CustomerId customerId, - ActionType actionType, Object... additionalInfo) { - String msgType = null; - switch (actionType) { - case ADDED: - msgType = DataConstants.ENTITY_CREATED; - break; - case DELETED: - msgType = DataConstants.ENTITY_DELETED; - break; - case UPDATED: - msgType = DataConstants.ENTITY_UPDATED; - break; - case ASSIGNED_TO_CUSTOMER: - msgType = DataConstants.ENTITY_ASSIGNED; - break; - case UNASSIGNED_FROM_CUSTOMER: - msgType = DataConstants.ENTITY_UNASSIGNED; - break; - case ATTRIBUTES_UPDATED: - msgType = DataConstants.ATTRIBUTES_UPDATED; - break; - case ATTRIBUTES_DELETED: - msgType = DataConstants.ATTRIBUTES_DELETED; - break; - case ALARM_ACK: - msgType = DataConstants.ALARM_ACK; - break; - case ALARM_CLEAR: - msgType = DataConstants.ALARM_CLEAR; - break; - case ASSIGNED_FROM_TENANT: - msgType = DataConstants.ENTITY_ASSIGNED_FROM_TENANT; - break; - case ASSIGNED_TO_TENANT: - msgType = DataConstants.ENTITY_ASSIGNED_TO_TENANT; - break; - case PROVISION_SUCCESS: - msgType = DataConstants.PROVISION_SUCCESS; - break; - case PROVISION_FAILURE: - msgType = DataConstants.PROVISION_FAILURE; - break; - case TIMESERIES_UPDATED: - msgType = DataConstants.TIMESERIES_UPDATED; - break; - case TIMESERIES_DELETED: - msgType = DataConstants.TIMESERIES_DELETED; - break; - case ASSIGNED_TO_EDGE: - msgType = DataConstants.ENTITY_ASSIGNED_TO_EDGE; - break; - case UNASSIGNED_FROM_EDGE: - msgType = DataConstants.ENTITY_UNASSIGNED_FROM_EDGE; - break; - } - if (!StringUtils.isEmpty(msgType)) { - try { - TbMsgMetaData metaData = new TbMsgMetaData(); - metaData.putValue("userId", user.getId().toString()); - metaData.putValue("userName", user.getName()); - if (customerId != null && !customerId.isNullUid()) { - metaData.putValue("customerId", customerId.toString()); - } - if (actionType == ActionType.ASSIGNED_TO_CUSTOMER) { - String strCustomerId = extractParameter(String.class, 1, additionalInfo); - String strCustomerName = extractParameter(String.class, 2, additionalInfo); - metaData.putValue("assignedCustomerId", strCustomerId); - metaData.putValue("assignedCustomerName", strCustomerName); - } else if (actionType == ActionType.UNASSIGNED_FROM_CUSTOMER) { - String strCustomerId = extractParameter(String.class, 1, additionalInfo); - String strCustomerName = extractParameter(String.class, 2, additionalInfo); - metaData.putValue("unassignedCustomerId", strCustomerId); - metaData.putValue("unassignedCustomerName", strCustomerName); - } else if (actionType == ActionType.ASSIGNED_FROM_TENANT) { - String strTenantId = extractParameter(String.class, 0, additionalInfo); - String strTenantName = extractParameter(String.class, 1, additionalInfo); - metaData.putValue("assignedFromTenantId", strTenantId); - metaData.putValue("assignedFromTenantName", strTenantName); - } else if (actionType == ActionType.ASSIGNED_TO_TENANT) { - String strTenantId = extractParameter(String.class, 0, additionalInfo); - String strTenantName = extractParameter(String.class, 1, additionalInfo); - metaData.putValue("assignedToTenantId", strTenantId); - metaData.putValue("assignedToTenantName", strTenantName); - } else if (actionType == ActionType.ASSIGNED_TO_EDGE) { - String strEdgeId = extractParameter(String.class, 1, additionalInfo); - String strEdgeName = extractParameter(String.class, 2, additionalInfo); - metaData.putValue("assignedEdgeId", strEdgeId); - metaData.putValue("assignedEdgeName", strEdgeName); - } else if (actionType == ActionType.UNASSIGNED_FROM_EDGE) { - String strEdgeId = extractParameter(String.class, 1, additionalInfo); - String strEdgeName = extractParameter(String.class, 2, additionalInfo); - metaData.putValue("unassignedEdgeId", strEdgeId); - metaData.putValue("unassignedEdgeName", strEdgeName); - } - ObjectNode entityNode; - if (entity != null) { - entityNode = json.valueToTree(entity); - if (entityId.getEntityType() == EntityType.DASHBOARD) { - entityNode.put("configuration", ""); - } - } else { - entityNode = json.createObjectNode(); - if (actionType == ActionType.ATTRIBUTES_UPDATED) { - String scope = extractParameter(String.class, 0, additionalInfo); - @SuppressWarnings("unchecked") - List attributes = extractParameter(List.class, 1, additionalInfo); - metaData.putValue(DataConstants.SCOPE, scope); - if (attributes != null) { - for (AttributeKvEntry attr : attributes) { - addKvEntry(entityNode, attr); - } - } - } else if (actionType == ActionType.ATTRIBUTES_DELETED) { - String scope = extractParameter(String.class, 0, additionalInfo); - @SuppressWarnings("unchecked") - List keys = extractParameter(List.class, 1, additionalInfo); - metaData.putValue(DataConstants.SCOPE, scope); - ArrayNode attrsArrayNode = entityNode.putArray("attributes"); - if (keys != null) { - keys.forEach(attrsArrayNode::add); - } - } else if (actionType == ActionType.TIMESERIES_UPDATED) { - @SuppressWarnings("unchecked") - List timeseries = extractParameter(List.class, 0, additionalInfo); - addTimeseries(entityNode, timeseries); - } else if (actionType == ActionType.TIMESERIES_DELETED) { - @SuppressWarnings("unchecked") - List keys = extractParameter(List.class, 0, additionalInfo); - if (keys != null) { - ArrayNode timeseriesArrayNode = entityNode.putArray("timeseries"); - keys.forEach(timeseriesArrayNode::add); - } - entityNode.put("startTs", extractParameter(Long.class, 1, additionalInfo)); - entityNode.put("endTs", extractParameter(Long.class, 2, additionalInfo)); - } - } - TbMsg tbMsg = TbMsg.newMsg(msgType, entityId, customerId, metaData, TbMsgDataType.JSON, json.writeValueAsString(entityNode)); - TenantId tenantId = user.getTenantId(); - if (tenantId.isNullUid()) { - if (entity instanceof HasTenantId) { - tenantId = ((HasTenantId) entity).getTenantId(); - } - } - tbClusterService.pushMsgToRuleEngine(tenantId, entityId, tbMsg, null); - } catch (Exception e) { - log.warn("[{}] Failed to push entity action to rule engine: {}", entityId, actionType, e); - } - } - } - - private void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) throws Exception { - if (kvEntry.getDataType() == DataType.BOOLEAN) { - kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); - } else if (kvEntry.getDataType() == DataType.DOUBLE) { - kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); - } else if (kvEntry.getDataType() == DataType.LONG) { - kvEntry.getLongValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); - } else if (kvEntry.getDataType() == DataType.JSON) { - if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(kvEntry.getKey(), json.readTree(kvEntry.getJsonValue().get())); - } - } else { - entityNode.put(kvEntry.getKey(), kvEntry.getValueAsString()); - } - } - - private T extractParameter(Class clazz, int index, Object... additionalInfo) { - T result = null; - if (additionalInfo != null && additionalInfo.length > index) { - Object paramObject = additionalInfo[index]; - if (clazz.isInstance(paramObject)) { - result = clazz.cast(paramObject); - } - } - return result; - } - protected String entityToStr(E entity) { try { return json.writeValueAsString(json.valueToTree(entity)); @@ -1095,23 +910,6 @@ public abstract class BaseController { return result; } - private void addTimeseries(ObjectNode entityNode, List timeseries) throws Exception { - if (timeseries != null && !timeseries.isEmpty()) { - ArrayNode result = entityNode.putArray("timeseries"); - Map> groupedTelemetry = timeseries.stream() - .collect(Collectors.groupingBy(TsKvEntry::getTs)); - for (Map.Entry> entry : groupedTelemetry.entrySet()) { - ObjectNode element = json.createObjectNode(); - element.put("ts", entry.getKey()); - ObjectNode values = element.putObject("values"); - for (TsKvEntry tsKvEntry : entry.getValue()) { - addKvEntry(values, tsKvEntry); - } - result.add(element); - } - } - } - protected void processDashboardIdFromAdditionalInfo(ObjectNode additionalInfo, String requiredFields) throws ThingsboardException { String dashboardId = additionalInfo.has(requiredFields) ? additionalInfo.get(requiredFields).asText() : null; if (dashboardId != null && !dashboardId.equals("null")) { diff --git a/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java b/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java new file mode 100644 index 0000000000..f1320d1a7a --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java @@ -0,0 +1,256 @@ +/** + * Copyright © 2016-2021 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.action; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.HasName; +import org.thingsboard.server.common.data.HasTenantId; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.audit.ActionType; +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.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.DataType; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgDataType; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.queue.TbClusterService; + +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +@TbCoreComponent +@Service +@RequiredArgsConstructor +@Slf4j +public class RuleEngineEntityActionService { + private final TbClusterService tbClusterService; + + private static final ObjectMapper json = new ObjectMapper(); + + public void pushEntityActionToRuleEngine(EntityId entityId, HasName entity, TenantId tenantId, CustomerId customerId, + ActionType actionType, User user, Object... additionalInfo) { + String msgType = null; + switch (actionType) { + case ADDED: + msgType = DataConstants.ENTITY_CREATED; + break; + case DELETED: + msgType = DataConstants.ENTITY_DELETED; + break; + case UPDATED: + msgType = DataConstants.ENTITY_UPDATED; + break; + case ASSIGNED_TO_CUSTOMER: + msgType = DataConstants.ENTITY_ASSIGNED; + break; + case UNASSIGNED_FROM_CUSTOMER: + msgType = DataConstants.ENTITY_UNASSIGNED; + break; + case ATTRIBUTES_UPDATED: + msgType = DataConstants.ATTRIBUTES_UPDATED; + break; + case ATTRIBUTES_DELETED: + msgType = DataConstants.ATTRIBUTES_DELETED; + break; + case ALARM_ACK: + msgType = DataConstants.ALARM_ACK; + break; + case ALARM_CLEAR: + msgType = DataConstants.ALARM_CLEAR; + break; + case ALARM_DELETE: + msgType = DataConstants.ALARM_DELETE; + break; + case ASSIGNED_FROM_TENANT: + msgType = DataConstants.ENTITY_ASSIGNED_FROM_TENANT; + break; + case ASSIGNED_TO_TENANT: + msgType = DataConstants.ENTITY_ASSIGNED_TO_TENANT; + break; + case PROVISION_SUCCESS: + msgType = DataConstants.PROVISION_SUCCESS; + break; + case PROVISION_FAILURE: + msgType = DataConstants.PROVISION_FAILURE; + break; + case TIMESERIES_UPDATED: + msgType = DataConstants.TIMESERIES_UPDATED; + break; + case TIMESERIES_DELETED: + msgType = DataConstants.TIMESERIES_DELETED; + break; + case ASSIGNED_TO_EDGE: + msgType = DataConstants.ENTITY_ASSIGNED_TO_EDGE; + break; + case UNASSIGNED_FROM_EDGE: + msgType = DataConstants.ENTITY_UNASSIGNED_FROM_EDGE; + break; + } + if (!StringUtils.isEmpty(msgType)) { + try { + TbMsgMetaData metaData = new TbMsgMetaData(); + if (user != null) { + metaData.putValue("userId", user.getId().toString()); + metaData.putValue("userName", user.getName()); + } + if (customerId != null && !customerId.isNullUid()) { + metaData.putValue("customerId", customerId.toString()); + } + if (actionType == ActionType.ASSIGNED_TO_CUSTOMER) { + String strCustomerId = extractParameter(String.class, 1, additionalInfo); + String strCustomerName = extractParameter(String.class, 2, additionalInfo); + metaData.putValue("assignedCustomerId", strCustomerId); + metaData.putValue("assignedCustomerName", strCustomerName); + } else if (actionType == ActionType.UNASSIGNED_FROM_CUSTOMER) { + String strCustomerId = extractParameter(String.class, 1, additionalInfo); + String strCustomerName = extractParameter(String.class, 2, additionalInfo); + metaData.putValue("unassignedCustomerId", strCustomerId); + metaData.putValue("unassignedCustomerName", strCustomerName); + } else if (actionType == ActionType.ASSIGNED_FROM_TENANT) { + String strTenantId = extractParameter(String.class, 0, additionalInfo); + String strTenantName = extractParameter(String.class, 1, additionalInfo); + metaData.putValue("assignedFromTenantId", strTenantId); + metaData.putValue("assignedFromTenantName", strTenantName); + } else if (actionType == ActionType.ASSIGNED_TO_TENANT) { + String strTenantId = extractParameter(String.class, 0, additionalInfo); + String strTenantName = extractParameter(String.class, 1, additionalInfo); + metaData.putValue("assignedToTenantId", strTenantId); + metaData.putValue("assignedToTenantName", strTenantName); + } else if (actionType == ActionType.ASSIGNED_TO_EDGE) { + String strEdgeId = extractParameter(String.class, 1, additionalInfo); + String strEdgeName = extractParameter(String.class, 2, additionalInfo); + metaData.putValue("assignedEdgeId", strEdgeId); + metaData.putValue("assignedEdgeName", strEdgeName); + } else if (actionType == ActionType.UNASSIGNED_FROM_EDGE) { + String strEdgeId = extractParameter(String.class, 1, additionalInfo); + String strEdgeName = extractParameter(String.class, 2, additionalInfo); + metaData.putValue("unassignedEdgeId", strEdgeId); + metaData.putValue("unassignedEdgeName", strEdgeName); + } + ObjectNode entityNode; + if (entity != null) { + entityNode = json.valueToTree(entity); + if (entityId.getEntityType() == EntityType.DASHBOARD) { + entityNode.put("configuration", ""); + } + } else { + entityNode = json.createObjectNode(); + if (actionType == ActionType.ATTRIBUTES_UPDATED) { + String scope = extractParameter(String.class, 0, additionalInfo); + @SuppressWarnings("unchecked") + List attributes = extractParameter(List.class, 1, additionalInfo); + metaData.putValue(DataConstants.SCOPE, scope); + if (attributes != null) { + for (AttributeKvEntry attr : attributes) { + addKvEntry(entityNode, attr); + } + } + } else if (actionType == ActionType.ATTRIBUTES_DELETED) { + String scope = extractParameter(String.class, 0, additionalInfo); + @SuppressWarnings("unchecked") + List keys = extractParameter(List.class, 1, additionalInfo); + metaData.putValue(DataConstants.SCOPE, scope); + ArrayNode attrsArrayNode = entityNode.putArray("attributes"); + if (keys != null) { + keys.forEach(attrsArrayNode::add); + } + } else if (actionType == ActionType.TIMESERIES_UPDATED) { + @SuppressWarnings("unchecked") + List timeseries = extractParameter(List.class, 0, additionalInfo); + addTimeseries(entityNode, timeseries); + } else if (actionType == ActionType.TIMESERIES_DELETED) { + @SuppressWarnings("unchecked") + List keys = extractParameter(List.class, 0, additionalInfo); + if (keys != null) { + ArrayNode timeseriesArrayNode = entityNode.putArray("timeseries"); + keys.forEach(timeseriesArrayNode::add); + } + entityNode.put("startTs", extractParameter(Long.class, 1, additionalInfo)); + entityNode.put("endTs", extractParameter(Long.class, 2, additionalInfo)); + } + } + TbMsg tbMsg = TbMsg.newMsg(msgType, entityId, customerId, metaData, TbMsgDataType.JSON, json.writeValueAsString(entityNode)); + if (tenantId.isNullUid()) { + if (entity instanceof HasTenantId) { + tenantId = ((HasTenantId) entity).getTenantId(); + } + } + tbClusterService.pushMsgToRuleEngine(tenantId, entityId, tbMsg, null); + } catch (Exception e) { + log.warn("[{}] Failed to push entity action to rule engine: {}", entityId, actionType, e); + } + } + } + + + private T extractParameter(Class clazz, int index, Object... additionalInfo) { + T result = null; + if (additionalInfo != null && additionalInfo.length > index) { + Object paramObject = additionalInfo[index]; + if (clazz.isInstance(paramObject)) { + result = clazz.cast(paramObject); + } + } + return result; + } + + private void addTimeseries(ObjectNode entityNode, List timeseries) throws Exception { + if (timeseries != null && !timeseries.isEmpty()) { + ArrayNode result = entityNode.putArray("timeseries"); + Map> groupedTelemetry = timeseries.stream() + .collect(Collectors.groupingBy(TsKvEntry::getTs)); + for (Map.Entry> entry : groupedTelemetry.entrySet()) { + ObjectNode element = json.createObjectNode(); + element.put("ts", entry.getKey()); + ObjectNode values = element.putObject("values"); + for (TsKvEntry tsKvEntry : entry.getValue()) { + addKvEntry(values, tsKvEntry); + } + result.add(element); + } + } + } + + private void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) throws Exception { + if (kvEntry.getDataType() == DataType.BOOLEAN) { + kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); + } else if (kvEntry.getDataType() == DataType.DOUBLE) { + kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); + } else if (kvEntry.getDataType() == DataType.LONG) { + kvEntry.getLongValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); + } else if (kvEntry.getDataType() == DataType.JSON) { + if (kvEntry.getJsonValue().isPresent()) { + entityNode.set(kvEntry.getKey(), json.readTree(kvEntry.getJsonValue().get())); + } + } else { + entityNode.put(kvEntry.getKey(), kvEntry.getValueAsString()); + } + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java index 95731d2988..05799fc643 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java @@ -17,9 +17,9 @@ package org.thingsboard.server.service.ttl; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; -import org.thingsboard.server.dao.util.PsqlDao; import java.sql.Connection; +import java.sql.DriverManager; import java.sql.ResultSet; import java.sql.SQLException; import java.sql.SQLWarning; @@ -62,4 +62,8 @@ public abstract class AbstractCleanUpService { protected abstract void doCleanUp(Connection connection) throws SQLException; + protected Connection getConnection() throws SQLException { + return DriverManager.getConnection(dbUrl, dbUserName, dbPassword); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java index 08b08a742b..7d6ebb6939 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java @@ -20,17 +20,27 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.audit.ActionType; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.dao.alarm.AlarmDao; +import org.thingsboard.server.dao.alarm.AlarmService; +import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantDao; import org.thingsboard.server.dao.util.PsqlDao; import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.service.action.RuleEngineEntityActionService; +import org.thingsboard.server.service.ttl.AbstractCleanUpService; +import java.sql.Connection; +import java.sql.SQLException; +import java.util.Date; import java.util.Optional; import java.util.UUID; import java.util.concurrent.TimeUnit; @@ -43,37 +53,54 @@ public class AlarmsCleanUpService { @Value("${sql.ttl.alarms.removal_batch_size}") private Integer removalBatchSize; - private final AlarmDao alarmDao; private final TenantDao tenantDao; + private final AlarmDao alarmDao; + private final AlarmService alarmService; + private final RelationService relationService; + private final RuleEngineEntityActionService ruleEngineEntityActionService; private final PartitionService partitionService; private final TbTenantProfileCache tenantProfileCache; - @Scheduled(initialDelayString = "${sql.ttl.alarms.checking_interval}", fixedDelayString = "${sql.ttl.alarms.checking_interval}") + @Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.alarms.checking_interval})}", fixedDelayString = "${sql.ttl.alarms.checking_interval}") public void cleanUp() { - PageLink tenantsBatchRequest = new PageLink(65536, 0); - PageLink alarmsRemovalBatchRequest = new PageLink(removalBatchSize, 0); - long currentTime = System.currentTimeMillis(); - + PageLink tenantsBatchRequest = new PageLink(10_000, 0); + PageLink removalBatchRequest = new PageLink(removalBatchSize, 0); PageData tenantsIds; do { tenantsIds = tenantDao.findTenantsIds(tenantsBatchRequest); - tenantsIds.getData().stream() - .filter(tenantId -> partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) - .forEach(tenantId -> { - Optional tenantProfileConfiguration = tenantProfileCache.get(tenantId).getProfileConfiguration(); - if (tenantProfileConfiguration.isEmpty() || tenantProfileConfiguration.get().getAlarmsTtlDays() == 0) { - return; - } + for (TenantId tenantId : tenantsIds.getData()) { + if (!partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) { + continue; + } + + Optional tenantProfileConfiguration = tenantProfileCache.get(tenantId).getProfileConfiguration(); + if (tenantProfileConfiguration.isEmpty() || tenantProfileConfiguration.get().getAlarmsTtlDays() == 0) { + continue; + } - PageData toRemove; - long outdatageTime = currentTime - TimeUnit.DAYS.toMillis(tenantProfileConfiguration.get().getAlarmsTtlDays()); - log.info("Cleaning up outdated alarms for tenant {}", tenantId); - do { - toRemove = alarmDao.findAlarmsIdsByEndTsBeforeAndTenantId(outdatageTime, tenantId, alarmsRemovalBatchRequest); - alarmDao.removeAllByIds(toRemove.getData()); - } while (toRemove.hasNext()); + long ttl = TimeUnit.DAYS.toMillis(tenantProfileConfiguration.get().getAlarmsTtlDays()); + long outdatageTime = System.currentTimeMillis() - ttl; + + long totalRemoved = 0; + while (true) { + PageData toRemove = alarmDao.findAlarmsIdsByEndTsBeforeAndTenantId(outdatageTime, tenantId, removalBatchRequest); + toRemove.getData().forEach(alarmId -> { + relationService.deleteEntityRelations(tenantId, alarmId); + Alarm alarm = alarmService.deleteAlarm(tenantId, alarmId).getAlarm(); + ruleEngineEntityActionService.pushEntityActionToRuleEngine(alarm.getOriginator(), alarm, tenantId, null, ActionType.ALARM_DELETE, null); }); + totalRemoved += toRemove.getTotalElements(); + if (!toRemove.hasNext()) { + break; + } + } + + if (totalRemoved > 0) { + log.info("Removed {} outdated alarm(s) for tenant {} older than {}", totalRemoved, tenantId, new Date(outdatageTime)); + } + } + tenantsBatchRequest = tenantsBatchRequest.nextPageLink(); } while (tenantsIds.hasNext()); } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java index 0c21719451..e93a82c7eb 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java @@ -40,7 +40,7 @@ public class EdgeEventsCleanUpService extends AbstractCleanUpService { @Scheduled(initialDelayString = "${sql.ttl.edge_events.execution_interval_ms}", fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}") public void cleanUp() { if (ttlTaskExecutionEnabled) { - try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { + try (Connection conn = getConnection()) { doCleanUp(conn); } catch (SQLException e) { log.error("SQLException occurred during TTL task execution ", e); diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java index 664e01e227..407c88261f 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java @@ -43,7 +43,7 @@ public class EventsCleanUpService extends AbstractCleanUpService { @Scheduled(initialDelayString = "${sql.ttl.events.execution_interval_ms}", fixedDelayString = "${sql.ttl.events.execution_interval_ms}") public void cleanUp() { if (ttlTaskExecutionEnabled) { - try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { + try (Connection conn = getConnection()) { doCleanUp(conn); } catch (SQLException e) { log.error("SQLException occurred during TTL task execution ", e); diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java index 9ece0b91a2..ee2d437a22 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java @@ -36,7 +36,7 @@ public abstract class AbstractTimeseriesCleanUpService extends AbstractCleanUpSe @Scheduled(initialDelayString = "${sql.ttl.ts.execution_interval_ms}", fixedDelayString = "${sql.ttl.ts.execution_interval_ms}") public void cleanUp() { if (ttlTaskExecutionEnabled) { - try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { + try (Connection conn = getConnection()) { doCleanUp(conn); } catch (SQLException e) { log.error("SQLException occurred during TTL task execution ", e); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 8edff5366e..8229046851 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -275,7 +275,7 @@ sql: edge_events_ttl: "${SQL_TTL_EDGE_EVENTS_TTL:2628000}" # Number of seconds. The current value corresponds to one month alarms: checking_interval: "${SQL_ALARMS_TTL_CHECKING_INTERVAL:7200000}" # Number of milliseconds. The current value corresponds to two hours - removal_batch_size: "${SQL_ALARMS_TTL_REMOVAL_BATCH_SIZE:200}" # To delete outdated alarms not all at once but in batches + removal_batch_size: "${SQL_ALARMS_TTL_REMOVAL_BATCH_SIZE:3000}" # To delete outdated alarms not all at once but in batches # Actor system parameters actors: diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 62459fc0ad..002cbbb733 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -66,6 +66,7 @@ public class DataConstants { public static final String TIMESERIES_DELETED = "TIMESERIES_DELETED"; public static final String ALARM_ACK = "ALARM_ACK"; public static final String ALARM_CLEAR = "ALARM_CLEAR"; + public static final String ALARM_DELETE = "ALARM_DELETE"; public static final String ENTITY_ASSIGNED_FROM_TENANT = "ENTITY_ASSIGNED_FROM_TENANT"; public static final String ENTITY_ASSIGNED_TO_TENANT = "ENTITY_ASSIGNED_TO_TENANT"; public static final String PROVISION_SUCCESS = "PROVISION_SUCCESS"; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/TenantProfile.java b/common/data/src/main/java/org/thingsboard/server/common/data/TenantProfile.java index e3bb6c6a4e..05d22af09f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/TenantProfile.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/TenantProfile.java @@ -93,6 +93,7 @@ public class TenantProfile extends SearchTextBased implements H } } + @JsonIgnore public Optional getProfileConfiguration() { return Optional.ofNullable(getProfileData().getConfiguration()) .filter(profileConfiguration -> profileConfiguration instanceof DefaultTenantProfileConfiguration) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java b/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java index 2594c73ebe..489c45f68d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java @@ -39,6 +39,7 @@ public enum ActionType { RELATIONS_DELETED(false), ALARM_ACK(false), ALARM_CLEAR(false), + ALARM_DELETE(false), LOGIN(false), LOGOUT(false), LOCKOUT(false), diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java index bcef402ff9..3b10f624d7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmQuery; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -56,6 +57,6 @@ public interface AlarmDao extends Dao { Set findAlarmSeverities(TenantId tenantId, EntityId entityId, Set status); - PageData findAlarmsIdsByEndTsBeforeAndTenantId(Long time, TenantId tenantId, PageLink pageLink); + PageData findAlarmsIdsByEndTsBeforeAndTenantId(Long time, TenantId tenantId, PageLink pageLink); } 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 b6eef91ac7..8f862390c3 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 @@ -23,6 +23,7 @@ import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.query.Param; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.model.sql.AlarmEntity; import org.thingsboard.server.dao.model.sql.AlarmInfoEntity; @@ -161,7 +162,7 @@ public interface AlarmRepository extends CrudRepository { @Param("affectedEntityType") String affectedEntityType, @Param("alarmStatuses") Set alarmStatuses); - @Query("SELECT a.id FROM AlarmEntity a WHERE a.createdTime < :time AND a.endTs < :time") - Page findAlarmsIdsByEndTsBefore(@Param("time") Long time, Pageable pageable); + @Query("SELECT a.id FROM AlarmEntity a WHERE a.tenantId = :tenantId AND a.createdTime < :time AND a.endTs < :time") + Page findAlarmsIdsByEndTsBeforeAndTenantId(@Param("time") Long time, @Param("tenantId") UUID tenantId, 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 f7da17d6ed..3216a23eee 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 @@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmQuery; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -164,7 +165,8 @@ public class JpaAlarmDao extends JpaAbstractDao implements A } @Override - public PageData findAlarmsIdsByEndTsBeforeAndTenantId(Long time, TenantId tenantId, PageLink pageLink) { - return DaoUtil.pageToPageData(alarmRepository.findAlarmsIdsByEndTsBefore(time, DaoUtil.toPageable(pageLink))); + public PageData findAlarmsIdsByEndTsBeforeAndTenantId(Long time, TenantId tenantId, PageLink pageLink) { + return DaoUtil.pageToPageData(alarmRepository.findAlarmsIdsByEndTsBeforeAndTenantId(time, tenantId.getId(), DaoUtil.toPageable(pageLink))) + .mapData(AlarmId::new); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java index e84beae702..69614f8079 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java @@ -150,6 +150,8 @@ class DeviceState { stateChanged = processAlarmClearNotification(ctx, msg); } else if (msg.getType().equals(DataConstants.ALARM_ACK)) { processAlarmAckNotification(ctx, msg); + } else if (msg.getType().equals(DataConstants.ALARM_DELETE)) { + processAlarmDeleteNotification(ctx, msg); } else { if (msg.getType().equals(DataConstants.ENTITY_ASSIGNED) || msg.getType().equals(DataConstants.ENTITY_UNASSIGNED)) { dynamicPredicateValueCtx.resetCustomer(); @@ -193,6 +195,12 @@ class DeviceState { ctx.tellSuccess(msg); } + private void processAlarmDeleteNotification(TbContext ctx, TbMsg msg) { + Alarm alarm = JacksonUtil.fromString(msg.getData(), Alarm.class); + alarmStates.values().removeIf(alarmState -> alarmState.getCurrentAlarm().getId().equals(alarm.getId())); + ctx.tellSuccess(msg); + } + private boolean processAttributesUpdateNotification(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException { String scope = msg.getMetaData().getValue(DataConstants.SCOPE); if (StringUtils.isEmpty(scope)) {