diff --git a/application/src/main/data/upgrade/3.2.2/schema_update.sql b/application/src/main/data/upgrade/3.2.2/schema_update.sql index 2814e18c2f..5647c301cf 100644 --- a/application/src/main/data/upgrade/3.2.2/schema_update.sql +++ b/application/src/main/data/upgrade/3.2.2/schema_update.sql @@ -67,6 +67,7 @@ CREATE TABLE IF NOT EXISTS ota_package ( type varchar(32) NOT NULL, title varchar(255) NOT NULL, version varchar(255) NOT NULL, + url varchar(255), file_name varchar(255), content_type varchar(255), checksum_algorithm varchar(32), diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index c54f940c38..b52f85af92 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -60,7 +60,9 @@ import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.event.EventService; import org.thingsboard.server.dao.nosql.CassandraBufferedRateExecutor; +import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.rule.RuleNodeStateService; import org.thingsboard.server.dao.tenant.TenantProfileService; @@ -311,6 +313,14 @@ public class ActorSystemContext { @Autowired(required = false) @Getter private EdgeRpcService edgeRpcService; + @Lazy + @Autowired(required = false) + @Getter private ResourceService resourceService; + + @Lazy + @Autowired(required = false) + @Getter private OtaPackageService otaPackageService; + @Value("${actors.session.max_concurrent_sessions_per_device:1}") @Getter private long maxConcurrentSessionsPerDevice; diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 1ddb99d7b7..9a1afb9ff8 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -69,7 +69,9 @@ import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.nosql.CassandraStatementTask; import org.thingsboard.server.dao.nosql.TbResultSetFuture; +import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -486,6 +488,16 @@ class DefaultTbContext implements TbContext { return mainCtx.getEntityViewService(); } + @Override + public ResourceService getResourceService() { + return mainCtx.getResourceService(); + } + + @Override + public OtaPackageService getOtaPackageService() { + return mainCtx.getOtaPackageService(); + } + @Override public RuleEngineDeviceProfileCache getDeviceProfileCache() { return mainCtx.getDeviceProfileCache(); 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 dd1d388a08..7d41deb2ab 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.ota.OtaPackageStateService; 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; @@ -279,6 +269,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; @@ -809,7 +802,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); } @@ -819,184 +812,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)); @@ -1093,23 +908,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/controller/OtaPackageController.java b/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java index 02e9d4b305..13d39b0b2a 100644 --- a/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java +++ b/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java @@ -64,6 +64,10 @@ public class OtaPackageController extends BaseController { OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); OtaPackage otaPackage = checkOtaPackageId(otaPackageId, Operation.READ); + if (otaPackage.hasUrl()) { + return ResponseEntity.badRequest().build(); + } + ByteArrayResource resource = new ByteArrayResource(otaPackage.getData().array()); return ResponseEntity.ok() .header(HttpHeaders.CONTENT_DISPOSITION, "attachment;filename=" + otaPackage.getFileName()) @@ -182,11 +186,10 @@ public class OtaPackageController extends BaseController { } @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") - @RequestMapping(value = "/otaPackages/{deviceProfileId}/{type}/{hasData}", method = RequestMethod.GET) + @RequestMapping(value = "/otaPackages/{deviceProfileId}/{type}", method = RequestMethod.GET) @ResponseBody public PageData getOtaPackages(@PathVariable("deviceProfileId") String strDeviceProfileId, @PathVariable("type") String strType, - @PathVariable("hasData") boolean hasData, @RequestParam int pageSize, @RequestParam int page, @RequestParam(required = false) String textSearch, @@ -197,7 +200,7 @@ public class OtaPackageController extends BaseController { try { PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); return checkNotNull(otaPackageService.findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(getTenantId(), - new DeviceProfileId(toUUID(strDeviceProfileId)), OtaPackageType.valueOf(strType), hasData, pageLink)); + new DeviceProfileId(toUUID(strDeviceProfileId)), OtaPackageType.valueOf(strType), pageLink)); } catch (Exception e) { throw handleException(e); } 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/ota/DefaultOtaPackageStateService.java b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java index c5d0c0472f..2ba4735d55 100644 --- a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java @@ -24,6 +24,7 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.OtaPackageInfo; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.id.TenantId; @@ -65,6 +66,7 @@ import static org.thingsboard.server.common.data.ota.OtaPackageKey.SIZE; import static org.thingsboard.server.common.data.ota.OtaPackageKey.STATE; import static org.thingsboard.server.common.data.ota.OtaPackageKey.TITLE; import static org.thingsboard.server.common.data.ota.OtaPackageKey.TS; +import static org.thingsboard.server.common.data.ota.OtaPackageKey.URL; import static org.thingsboard.server.common.data.ota.OtaPackageKey.VERSION; import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE; import static org.thingsboard.server.common.data.ota.OtaPackageType.SOFTWARE; @@ -261,11 +263,12 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService { } - private void update(Device device, OtaPackageInfo firmware, long ts) { + private void update(Device device, OtaPackageInfo otaPackage, long ts) { TenantId tenantId = device.getTenantId(); DeviceId deviceId = device.getId(); + OtaPackageType otaPackageType = otaPackage.getType(); - BasicTsKvEntry status = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(getTelemetryKey(firmware.getType(), STATE), OtaPackageUpdateStatus.INITIATED.name())); + BasicTsKvEntry status = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(getTelemetryKey(otaPackageType, STATE), OtaPackageUpdateStatus.INITIATED.name())); telemetryService.saveAndNotify(tenantId, deviceId, Collections.singletonList(status), new FutureCallback<>() { @Override @@ -280,11 +283,37 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService { }); List attributes = new ArrayList<>(); - attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(firmware.getType(), TITLE), firmware.getTitle()))); - attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(firmware.getType(), VERSION), firmware.getVersion()))); - attributes.add(new BaseAttributeKvEntry(ts, new LongDataEntry(getAttributeKey(firmware.getType(), SIZE), firmware.getDataSize()))); - attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(firmware.getType(), CHECKSUM_ALGORITHM), firmware.getChecksumAlgorithm().name()))); - attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(firmware.getType(), CHECKSUM), firmware.getChecksum()))); + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, TITLE), otaPackage.getTitle()))); + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, VERSION), otaPackage.getVersion()))); + if (otaPackage.hasUrl()) { + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, URL), otaPackage.getUrl()))); + List attrToRemove = new ArrayList<>(); + + if (otaPackage.getDataSize() == null) { + attrToRemove.add(getAttributeKey(otaPackageType, SIZE)); + } else { + attributes.add(new BaseAttributeKvEntry(ts, new LongDataEntry(getAttributeKey(otaPackageType, SIZE), otaPackage.getDataSize()))); + } + + if (otaPackage.getChecksumAlgorithm() != null) { + attrToRemove.add(getAttributeKey(otaPackageType, CHECKSUM_ALGORITHM)); + } else { + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM_ALGORITHM), otaPackage.getChecksumAlgorithm().name()))); + } + + if (StringUtils.isEmpty(otaPackage.getChecksum())) { + attrToRemove.add(getAttributeKey(otaPackageType, CHECKSUM)); + } else { + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM), otaPackage.getChecksum()))); + } + + remove(device, otaPackageType, attrToRemove); + } else { + attributes.add(new BaseAttributeKvEntry(ts, new LongDataEntry(getAttributeKey(otaPackageType, SIZE), otaPackage.getDataSize()))); + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM_ALGORITHM), otaPackage.getChecksumAlgorithm().name()))); + attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM), otaPackage.getChecksum()))); + remove(device, otaPackageType, Collections.singletonList(getAttributeKey(otaPackageType, URL))); + } telemetryService.saveAndNotify(tenantId, deviceId, DataConstants.SHARED_SCOPE, attributes, new FutureCallback<>() { @Override @@ -299,20 +328,24 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService { }); } - private void remove(Device device, OtaPackageType firmwareType) { - telemetryService.deleteAndNotify(device.getTenantId(), device.getId(), DataConstants.SHARED_SCOPE, OtaPackageUtil.getAttributeKeys(firmwareType), + private void remove(Device device, OtaPackageType otaPackageType) { + remove(device, otaPackageType, OtaPackageUtil.getAttributeKeys(otaPackageType)); + } + + private void remove(Device device, OtaPackageType otaPackageType, List attributesKeys) { + telemetryService.deleteAndNotify(device.getTenantId(), device.getId(), DataConstants.SHARED_SCOPE, attributesKeys, new FutureCallback<>() { @Override public void onSuccess(@Nullable Void tmp) { - log.trace("[{}] Success remove target firmware attributes!", device.getId()); + log.trace("[{}] Success remove target {} attributes!", device.getId(), otaPackageType); Set keysToNotify = new HashSet<>(); - OtaPackageUtil.ALL_FW_ATTRIBUTE_KEYS.forEach(key -> keysToNotify.add(new AttributeKey(DataConstants.SHARED_SCOPE, key))); + attributesKeys.forEach(key -> keysToNotify.add(new AttributeKey(DataConstants.SHARED_SCOPE, key))); tbClusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete(device.getTenantId(), device.getId(), keysToNotify), null); } @Override public void onFailure(Throwable t) { - log.error("[{}] Failed to remove target firmware attributes!", device.getId(), t); + log.error("[{}] Failed to remove target {} attributes!", device.getId(), otaPackageType, t); } }); } diff --git a/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java b/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java index 2cde0e2113..2c3e7f5200 100644 --- a/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java +++ b/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java @@ -157,6 +157,11 @@ public class DefaultTbResourceService implements TbResourceService { resourceService.deleteResourcesByTenantId(tenantId); } + @Override + public long sumDataSizeByTenantId(TenantId tenantId) { + return resourceService.sumDataSizeByTenantId(tenantId); + } + private Comparator getComparator(String sortProperty, String sortOrder) { Comparator comparator; if ("name".equals(sortProperty)) { diff --git a/application/src/main/java/org/thingsboard/server/service/resource/TbResourceService.java b/application/src/main/java/org/thingsboard/server/service/resource/TbResourceService.java index 7ad1848138..d3d079f548 100644 --- a/application/src/main/java/org/thingsboard/server/service/resource/TbResourceService.java +++ b/application/src/main/java/org/thingsboard/server/service/resource/TbResourceService.java @@ -55,4 +55,5 @@ public interface TbResourceService { void deleteResourcesByTenantId(TenantId tenantId); + long sumDataSizeByTenantId(TenantId tenantId); } diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 76bbe1f518..980c3d2c10 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -536,6 +536,9 @@ public class DefaultTransportApiService implements TransportApiService { if (otaPackageInfo == null) { builder.setResponseStatus(TransportProtos.ResponseStatus.NOT_FOUND); + } else if (otaPackageInfo.hasUrl()) { + builder.setResponseStatus(TransportProtos.ResponseStatus.FAILURE); + log.trace("[{}] Can`t send OtaPackage with URL data!", otaPackageInfo.getId()); } else { builder.setResponseStatus(TransportProtos.ResponseStatus.SUCCESS); builder.setOtaPackageIdMSB(otaPackageId.getId().getMostSignificantBits()); 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 new file mode 100644 index 0000000000..3b76a6cbca --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java @@ -0,0 +1,110 @@ +/** + * 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.ttl.alarms; + +import lombok.RequiredArgsConstructor; +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.page.SortOrder; +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.queue.util.TbCoreComponent; +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; + +@TbCoreComponent +@Service +@Slf4j +@RequiredArgsConstructor +public class AlarmsCleanUpService { + @Value("${sql.ttl.alarms.removal_batch_size}") + private Integer removalBatchSize; + + 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 = "#{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(10_000, 0); + PageLink removalBatchRequest = new PageLink(removalBatchSize, 0 ); + PageData tenantsIds; + do { + tenantsIds = tenantDao.findTenantsIds(tenantsBatchRequest); + 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; + } + + long ttl = TimeUnit.DAYS.toMillis(tenantProfileConfiguration.get().getAlarmsTtlDays()); + long expirationTime = System.currentTimeMillis() - ttl; + + long totalRemoved = 0; + while (true) { + PageData toRemove = alarmDao.findAlarmsIdsByEndTsBeforeAndTenantId(expirationTime, 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(expirationTime)); + } + } + + 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 de1f3cd3af..78658fe8fd 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -273,6 +273,9 @@ sql: enabled: "${SQL_TTL_EDGE_EVENTS_ENABLED:true}" execution_interval_ms: "${SQL_TTL_EDGE_EVENTS_EXECUTION_INTERVAL:86400000}" # Number of milliseconds. The current value corresponds to one day 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:3000}" # To delete outdated alarms not all at once but in batches # Actor system parameters actors: diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseOtaPackageControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseOtaPackageControllerTest.java index 9b73053831..a11d1b624a 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseOtaPackageControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseOtaPackageControllerTest.java @@ -50,7 +50,7 @@ public abstract class BaseOtaPackageControllerTest extends AbstractControllerTes private static final String FILE_NAME = "filename.txt"; private static final String VERSION = "v1.0"; private static final String CONTENT_TYPE = "text/plain"; - private static final String CHECKSUM_ALGORITHM = "sha256"; + private static final String CHECKSUM_ALGORITHM = "SHA256"; private static final String CHECKSUM = "4bf5122f344554c53bde2ebb8cd2b7e3d1600ad631c385a5d7cce23c7785459a"; private static final ByteBuffer DATA = ByteBuffer.wrap(new byte[]{1}); @@ -257,7 +257,7 @@ public abstract class BaseOtaPackageControllerTest extends AbstractControllerTes @Test public void testFindTenantFirmwaresByHasData() throws Exception { List otaPackagesWithData = new ArrayList<>(); - List otaPackagesWithoutData = new ArrayList<>(); + List allOtaPackages = new ArrayList<>(); for (int i = 0; i < 165; i++) { OtaPackageInfo firmwareInfo = new OtaPackageInfo(); @@ -272,44 +272,45 @@ public abstract class BaseOtaPackageControllerTest extends AbstractControllerTes MockMultipartFile testData = new MockMultipartFile("file", FILE_NAME, CONTENT_TYPE, DATA.array()); OtaPackage savedFirmware = savaData("/api/otaPackage/" + savedFirmwareInfo.getId().getId().toString() + "?checksum={checksum}&checksumAlgorithm={checksumAlgorithm}", testData, CHECKSUM, CHECKSUM_ALGORITHM); - otaPackagesWithData.add(new OtaPackageInfo(savedFirmware)); - } else { - otaPackagesWithoutData.add(savedFirmwareInfo); + savedFirmwareInfo = new OtaPackageInfo(savedFirmware); + otaPackagesWithData.add(savedFirmwareInfo); } + + allOtaPackages.add(savedFirmwareInfo); } - List loadedFirmwaresWithData = new ArrayList<>(); + List loadedOtaPackagesWithData = new ArrayList<>(); PageLink pageLink = new PageLink(24); PageData pageData; do { - pageData = doGetTypedWithPageLink("/api/otaPackages/" + deviceProfileId.toString() + "/FIRMWARE/true?", + pageData = doGetTypedWithPageLink("/api/otaPackages/" + deviceProfileId.toString() + "/FIRMWARE?", new TypeReference<>() { }, pageLink); - loadedFirmwaresWithData.addAll(pageData.getData()); + loadedOtaPackagesWithData.addAll(pageData.getData()); if (pageData.hasNext()) { pageLink = pageLink.nextPageLink(); } } while (pageData.hasNext()); - List loadedFirmwaresWithoutData = new ArrayList<>(); + List allLoadedOtaPackages = new ArrayList<>(); pageLink = new PageLink(24); do { - pageData = doGetTypedWithPageLink("/api/otaPackages/" + deviceProfileId.toString() + "/FIRMWARE/false?", + pageData = doGetTypedWithPageLink("/api/otaPackages?", new TypeReference<>() { }, pageLink); - loadedFirmwaresWithoutData.addAll(pageData.getData()); + allLoadedOtaPackages.addAll(pageData.getData()); if (pageData.hasNext()) { pageLink = pageLink.nextPageLink(); } } while (pageData.hasNext()); Collections.sort(otaPackagesWithData, idComparator); - Collections.sort(otaPackagesWithoutData, idComparator); - Collections.sort(loadedFirmwaresWithData, idComparator); - Collections.sort(loadedFirmwaresWithoutData, idComparator); + Collections.sort(allOtaPackages, idComparator); + Collections.sort(loadedOtaPackagesWithData, idComparator); + Collections.sort(allLoadedOtaPackages, idComparator); - Assert.assertEquals(otaPackagesWithData, loadedFirmwaresWithData); - Assert.assertEquals(otaPackagesWithoutData, loadedFirmwaresWithoutData); + Assert.assertEquals(otaPackagesWithData, loadedOtaPackagesWithData); + Assert.assertEquals(allOtaPackages, allLoadedOtaPackages); } diff --git a/application/src/test/java/org/thingsboard/server/service/resource/BaseTbResourceServiceTest.java b/application/src/test/java/org/thingsboard/server/service/resource/BaseTbResourceServiceTest.java index 464ca5c3cf..62facbb424 100644 --- a/application/src/test/java/org/thingsboard/server/service/resource/BaseTbResourceServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/resource/BaseTbResourceServiceTest.java @@ -19,17 +19,24 @@ import com.datastax.oss.driver.api.core.uuid.Uuids; import org.junit.After; import org.junit.Assert; import org.junit.Before; +import org.junit.Rule; import org.junit.Test; +import org.junit.rules.ExpectedException; import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.common.data.EntityInfo; +import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.ResourceType; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.exception.ThingsboardException; 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.security.Authority; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DaoSqlTest; @@ -109,6 +116,64 @@ public class BaseTbResourceServiceTest extends AbstractControllerTest { .andExpect(status().isOk()); } + @Rule + public ExpectedException thrown = ExpectedException.none(); + + @Test + public void testSaveResourceWithMaxSumDataSizeOutOfLimit() throws Exception { + loginSysAdmin(); + long limit = 1; + EntityInfo defaultTenantProfileInfo = doGet("/api/tenantProfileInfo/default", EntityInfo.class); + TenantProfile defaultTenantProfile = doGet("/api/tenantProfile/" + defaultTenantProfileInfo.getId().getId().toString(), TenantProfile.class); + defaultTenantProfile.getProfileData().setConfiguration(DefaultTenantProfileConfiguration.builder().maxResourcesInBytes(limit).build()); + doPost("/api/tenantProfile", defaultTenantProfile, TenantProfile.class); + + loginTenantAdmin(); + + Assert.assertEquals(0, resourceService.sumDataSizeByTenantId(tenantId)); + + createResource("test", DEFAULT_FILE_NAME); + + Assert.assertEquals(1, resourceService.sumDataSizeByTenantId(tenantId)); + + try { + thrown.expect(DataValidationException.class); + thrown.expectMessage(String.format("Failed to create the tb resource, files size limit is exhausted %d bytes!", limit)); + createResource("test1", 1 + DEFAULT_FILE_NAME); + } finally { + defaultTenantProfile.getProfileData().setConfiguration(DefaultTenantProfileConfiguration.builder().maxResourcesInBytes(0).build()); + loginSysAdmin(); + doPost("/api/tenantProfile", defaultTenantProfile, TenantProfile.class); + } + } + + @Test + public void sumDataSizeByTenantId() throws ThingsboardException { + Assert.assertEquals(0, resourceService.sumDataSizeByTenantId(tenantId)); + + createResource("test", DEFAULT_FILE_NAME); + Assert.assertEquals(1, resourceService.sumDataSizeByTenantId(tenantId)); + + int maxSumDataSize = 8; + + for (int i = 2; i <= maxSumDataSize; i++) { + createResource("test" + i, i + DEFAULT_FILE_NAME); + Assert.assertEquals(i, resourceService.sumDataSizeByTenantId(tenantId)); + } + + Assert.assertEquals(maxSumDataSize, resourceService.sumDataSizeByTenantId(tenantId)); + } + + private TbResource createResource(String title, String filename) throws ThingsboardException { + TbResource resource = new TbResource(); + resource.setTenantId(tenantId); + resource.setTitle(title); + resource.setResourceType(ResourceType.JKS); + resource.setFileName(filename); + resource.setData("1"); + return resourceService.saveResource(resource); + } + @Test public void testSaveTbResource() throws Exception { TbResource resource = new TbResource(); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/OtaPackageService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/OtaPackageService.java index 589bdf14b6..fea29681c1 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/OtaPackageService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/OtaPackageService.java @@ -44,9 +44,11 @@ public interface OtaPackageService { PageData findTenantOtaPackagesByTenantId(TenantId tenantId, PageLink pageLink); - PageData findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, boolean hasData, PageLink pageLink); + PageData findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, PageLink pageLink); void deleteOtaPackage(TenantId tenantId, OtaPackageId otaPackageId); void deleteOtaPackagesByTenantId(TenantId tenantId); + + long sumDataSizeByTenantId(TenantId tenantId); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java index 802628ff2c..1694c89bae 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java @@ -49,5 +49,5 @@ public interface ResourceService { void deleteResourcesByTenantId(TenantId tenantId); - + long sumDataSizeByTenantId(TenantId tenantId); } 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/OtaPackage.java b/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackage.java index 6110310cd3..3506aaea75 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackage.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackage.java @@ -37,8 +37,8 @@ public class OtaPackage extends OtaPackageInfo { super(id); } - public OtaPackage(OtaPackage firmware) { - super(firmware); - this.data = firmware.getData(); + public OtaPackage(OtaPackage otaPackage) { + super(otaPackage); + this.data = otaPackage.getData(); } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackageInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackageInfo.java index 5a33a95215..f27c90c20b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackageInfo.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/OtaPackageInfo.java @@ -37,6 +37,7 @@ public class OtaPackageInfo extends SearchTextBasedWithAdditionalInfo implements H } } + @JsonIgnore + public Optional getProfileConfiguration() { + return Optional.ofNullable(getProfileData().getConfiguration()) + .filter(profileConfiguration -> profileConfiguration instanceof DefaultTenantProfileConfiguration) + .map(profileConfiguration -> (DefaultTenantProfileConfiguration) profileConfiguration); + } + public TenantProfileData createDefaultTenantProfileData() { TenantProfileData tpd = new TenantProfileData(); tpd.setConfiguration(new 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/common/data/src/main/java/org/thingsboard/server/common/data/exception/ApiUsageLimitsExceededException.java b/common/data/src/main/java/org/thingsboard/server/common/data/exception/ApiUsageLimitsExceededException.java new file mode 100644 index 0000000000..fb5c4dffeb --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/exception/ApiUsageLimitsExceededException.java @@ -0,0 +1,25 @@ +/** + * 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.common.data.exception; + +public class ApiUsageLimitsExceededException extends RuntimeException { + public ApiUsageLimitsExceededException(String message) { + super(message); + } + + public ApiUsageLimitsExceededException() { + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ota/OtaPackageKey.java b/common/data/src/main/java/org/thingsboard/server/common/data/ota/OtaPackageKey.java index 0528b9dfe3..8b9cdfa8fb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ota/OtaPackageKey.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ota/OtaPackageKey.java @@ -19,7 +19,7 @@ import lombok.Getter; public enum OtaPackageKey { - TITLE("title"), VERSION("version"), TS("ts"), STATE("state"), SIZE("size"), CHECKSUM("checksum"), CHECKSUM_ALGORITHM("checksum_algorithm"); + TITLE("title"), VERSION("version"), TS("ts"), STATE("state"), SIZE("size"), CHECKSUM("checksum"), CHECKSUM_ALGORITHM("checksum_algorithm"), URL("url"); @Getter private final String value; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java index 2020245ea1..6ffbce4d3d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java @@ -17,10 +17,11 @@ package org.thingsboard.server.common.data.page; import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonProperty; -import org.thingsboard.server.common.data.BaseData; import java.util.Collections; import java.util.List; +import java.util.function.Function; +import java.util.stream.Collectors; public class PageData { @@ -61,4 +62,8 @@ public class PageData { return hasNext; } + public PageData mapData(Function mapper) { + return new PageData<>(getData().stream().map(mapper).collect(Collectors.toList()), getTotalPages(), getTotalElements(), hasNext()); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index b9bd72b0db..8cdccfe8bd 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -34,6 +34,8 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxUsers; private long maxDashboards; private long maxRuleChains; + private long maxResourcesInBytes; + private long maxOtaPackagesInBytes; private String transportTenantMsgRateLimit; private String transportTenantTelemetryMsgRateLimit; @@ -53,6 +55,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxCreatedAlarms; private int defaultStorageTtlDays; + private int alarmsTtlDays; private double warnThreshold; diff --git a/dao/src/main/java/org/thingsboard/server/dao/Dao.java b/dao/src/main/java/org/thingsboard/server/dao/Dao.java index 0111abdbbe..a5d4dfd9d1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/Dao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/Dao.java @@ -18,6 +18,7 @@ package org.thingsboard.server.dao; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.id.TenantId; +import java.util.Collection; import java.util.List; import java.util.UUID; @@ -33,4 +34,6 @@ public interface Dao { boolean removeById(TenantId tenantId, UUID id); + void removeAllByIds(Collection ids); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/TenantEntityWithDataDao.java b/dao/src/main/java/org/thingsboard/server/dao/TenantEntityWithDataDao.java new file mode 100644 index 0000000000..2199a51cd4 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/TenantEntityWithDataDao.java @@ -0,0 +1,23 @@ +/** + * 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.dao; + +import org.thingsboard.server.common.data.id.TenantId; + +public interface TenantEntityWithDataDao { + + Long sumDataSizeByTenantId(TenantId tenantId); +} 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 eb873db679..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,10 +21,12 @@ 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; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.dao.Dao; @@ -54,4 +56,7 @@ public interface AlarmDao extends Dao { AlarmDataQuery query, Collection orderedEntityIds); Set findAlarmSeverities(TenantId tenantId, EntityId entityId, Set status); + + PageData findAlarmsIdsByEndTsBeforeAndTenantId(Long time, TenantId tenantId, PageLink pageLink); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java index bb54efe769..90d988c20b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java @@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.alarm.AlarmQuery; import org.thingsboard.server.common.data.alarm.AlarmSearchStatus; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.exception.ApiUsageLimitsExceededException; import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; @@ -119,7 +120,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ Alarm existing = alarmDao.findLatestByOriginatorAndType(alarm.getTenantId(), alarm.getOriginator(), alarm.getType()).get(); if (existing == null || existing.getStatus().isCleared()) { if (!alarmCreationEnabled) { - throw new IllegalStateException("Alarm creation is disabled"); + throw new ApiUsageLimitsExceededException("Alarms creation is disabled"); } return createAlarm(alarm); } else { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 34878a3810..bfb530dd35 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -487,6 +487,7 @@ public class ModelConstants { public static final String OTA_PACKAGE_TYPE_COLUMN = "type"; public static final String OTA_PACKAGE_TILE_COLUMN = TITLE_PROPERTY; public static final String OTA_PACKAGE_VERSION_COLUMN = "version"; + public static final String OTA_PACKAGE_URL_COLUMN = "url"; public static final String OTA_PACKAGE_FILE_NAME_COLUMN = "file_name"; public static final String OTA_PACKAGE_CONTENT_TYPE_COLUMN = "content_type"; public static final String OTA_PACKAGE_CHECKSUM_ALGORITHM_COLUMN = "checksum_algorithm"; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageEntity.java index 97e4dbebbd..a5291e8a79 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageEntity.java @@ -51,6 +51,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TABLE_ import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TENANT_ID_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TILE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TYPE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_URL_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_VERSION_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.SEARCH_TEXT_PROPERTY; @@ -77,6 +78,9 @@ public class OtaPackageEntity extends BaseSqlEntity implements Searc @Column(name = OTA_PACKAGE_VERSION_COLUMN) private String version; + @Column(name = OTA_PACKAGE_URL_COLUMN) + private String url; + @Column(name = OTA_PACKAGE_FILE_NAME_COLUMN) private String fileName; @@ -118,6 +122,7 @@ public class OtaPackageEntity extends BaseSqlEntity implements Searc this.type = firmware.getType(); this.title = firmware.getTitle(); this.version = firmware.getVersion(); + this.url = firmware.getUrl(); this.fileName = firmware.getFileName(); this.contentType = firmware.getContentType(); this.checksumAlgorithm = firmware.getChecksumAlgorithm(); @@ -148,6 +153,7 @@ public class OtaPackageEntity extends BaseSqlEntity implements Searc firmware.setType(type); firmware.setTitle(title); firmware.setVersion(version); + firmware.setUrl(url); firmware.setFileName(fileName); firmware.setContentType(contentType); firmware.setChecksumAlgorithm(checksumAlgorithm); diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageInfoEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageInfoEntity.java index 30441ed098..db16251f71 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageInfoEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/OtaPackageInfoEntity.java @@ -22,6 +22,7 @@ import org.hibernate.annotations.Type; import org.hibernate.annotations.TypeDef; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.OtaPackageInfo; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.id.DeviceProfileId; @@ -50,6 +51,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TABLE_ import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TENANT_ID_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TILE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_TYPE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_URL_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.OTA_PACKAGE_VERSION_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.SEARCH_TEXT_PROPERTY; @@ -76,6 +78,9 @@ public class OtaPackageInfoEntity extends BaseSqlEntity implemen @Column(name = OTA_PACKAGE_VERSION_COLUMN) private String version; + @Column(name = OTA_PACKAGE_URL_COLUMN) + private String url; + @Column(name = OTA_PACKAGE_FILE_NAME_COLUMN) private String fileName; @@ -116,6 +121,7 @@ public class OtaPackageInfoEntity extends BaseSqlEntity implemen } this.title = firmware.getTitle(); this.version = firmware.getVersion(); + this.url = firmware.getUrl(); this.fileName = firmware.getFileName(); this.contentType = firmware.getContentType(); this.checksumAlgorithm = firmware.getChecksumAlgorithm(); @@ -125,7 +131,7 @@ public class OtaPackageInfoEntity extends BaseSqlEntity implemen } public OtaPackageInfoEntity(UUID id, long createdTime, UUID tenantId, UUID deviceProfileId, OtaPackageType type, String title, String version, - String fileName, String contentType, ChecksumAlgorithm checksumAlgorithm, String checksum, Long dataSize, + String url, String fileName, String contentType, ChecksumAlgorithm checksumAlgorithm, String checksum, Long dataSize, Object additionalInfo, boolean hasData) { this.id = id; this.createdTime = createdTime; @@ -134,6 +140,7 @@ public class OtaPackageInfoEntity extends BaseSqlEntity implemen this.type = type; this.title = title; this.version = version; + this.url = url; this.fileName = fileName; this.contentType = contentType; this.checksumAlgorithm = checksumAlgorithm; @@ -164,6 +171,7 @@ public class OtaPackageInfoEntity extends BaseSqlEntity implemen firmware.setType(type); firmware.setTitle(title); firmware.setVersion(version); + firmware.setUrl(url); firmware.setFileName(fileName); firmware.setContentType(contentType); firmware.setChecksumAlgorithm(checksumAlgorithm); diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java index 536c79843a..881baea4f1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java @@ -20,28 +20,32 @@ import com.google.common.hash.Hashing; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; import org.hibernate.exception.ConstraintViolationException; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.Cache; import org.springframework.cache.CacheManager; import org.springframework.cache.annotation.Cacheable; +import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.server.cache.ota.OtaPackageDataCache; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.OtaPackageInfo; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.Tenant; -import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; -import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; +import org.thingsboard.server.common.data.ota.OtaPackageType; 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.dao.device.DeviceProfileDao; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantDao; import java.nio.ByteBuffer; @@ -50,6 +54,7 @@ import java.util.List; import java.util.Optional; import static org.thingsboard.server.common.data.CacheConstants.OTA_PACKAGE_CACHE; +import static org.thingsboard.server.common.data.EntityType.OTA_PACKAGE; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @@ -67,6 +72,10 @@ public class BaseOtaPackageService implements OtaPackageService { private final CacheManager cacheManager; private final OtaPackageDataCache otaPackageDataCache; + @Autowired + @Lazy + private TbTenantProfileCache tenantProfileCache; + @Override public OtaPackageInfo saveOtaPackageInfo(OtaPackageInfo otaPackageInfo) { log.trace("Executing saveOtaPackageInfo [{}]", otaPackageInfo); @@ -172,11 +181,11 @@ public class BaseOtaPackageService implements OtaPackageService { } @Override - public PageData findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, boolean hasData, PageLink pageLink) { - log.trace("Executing findTenantOtaPackagesByTenantIdAndHasData, tenantId [{}], hasData [{}] pageLink [{}]", tenantId, hasData, pageLink); + public PageData findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, PageLink pageLink) { + log.trace("Executing findTenantOtaPackagesByTenantIdAndHasData, tenantId [{}], pageLink [{}]", tenantId, pageLink); validateId(tenantId, INCORRECT_TENANT_ID + tenantId); validatePageLink(pageLink); - return otaPackageInfoDao.findOtaPackageInfoByTenantIdAndDeviceProfileIdAndTypeAndHasData(tenantId, deviceProfileId, otaPackageType, hasData, pageLink); + return otaPackageInfoDao.findOtaPackageInfoByTenantIdAndDeviceProfileIdAndTypeAndHasData(tenantId, deviceProfileId, otaPackageType, pageLink); } @Override @@ -204,6 +213,11 @@ public class BaseOtaPackageService implements OtaPackageService { } } + @Override + public long sumDataSizeByTenantId(TenantId tenantId) { + return otaPackageDao.sumDataSizeByTenantId(tenantId); + } + @Override public void deleteOtaPackagesByTenantId(TenantId tenantId) { log.trace("Executing deleteOtaPackagesByTenantId, tenantId [{}]", tenantId); @@ -227,31 +241,43 @@ public class BaseOtaPackageService implements OtaPackageService { private DataValidator otaPackageValidator = new DataValidator<>() { + @Override + protected void validateCreate(TenantId tenantId, OtaPackage otaPackage) { + DefaultTenantProfileConfiguration profileConfiguration = + (DefaultTenantProfileConfiguration) tenantProfileCache.get(tenantId).getProfileData().getConfiguration(); + long maxOtaPackagesInBytes = profileConfiguration.getMaxOtaPackagesInBytes(); + validateMaxSumDataSizePerTenant(tenantId, otaPackageDao, maxOtaPackagesInBytes, otaPackage.getDataSize(), OTA_PACKAGE); + } + @Override protected void validateDataImpl(TenantId tenantId, OtaPackage otaPackage) { validateImpl(otaPackage); - if (StringUtils.isEmpty(otaPackage.getFileName())) { - throw new DataValidationException("OtaPackage file name should be specified!"); - } + if (!otaPackage.hasUrl()) { + if (StringUtils.isEmpty(otaPackage.getFileName())) { + throw new DataValidationException("OtaPackage file name should be specified!"); + } - if (StringUtils.isEmpty(otaPackage.getContentType())) { - throw new DataValidationException("OtaPackage content type should be specified!"); - } + if (StringUtils.isEmpty(otaPackage.getContentType())) { + throw new DataValidationException("OtaPackage content type should be specified!"); + } - if (otaPackage.getChecksumAlgorithm() == null) { - throw new DataValidationException("OtaPackage checksum algorithm should be specified!"); - } - if (StringUtils.isEmpty(otaPackage.getChecksum())) { - throw new DataValidationException("OtaPackage checksum should be specified!"); - } + if (otaPackage.getChecksumAlgorithm() == null) { + throw new DataValidationException("OtaPackage checksum algorithm should be specified!"); + } + if (StringUtils.isEmpty(otaPackage.getChecksum())) { + throw new DataValidationException("OtaPackage checksum should be specified!"); + } - String currentChecksum; + String currentChecksum; - currentChecksum = generateChecksum(otaPackage.getChecksumAlgorithm(), otaPackage.getData()); + currentChecksum = generateChecksum(otaPackage.getChecksumAlgorithm(), otaPackage.getData()); - if (!currentChecksum.equals(otaPackage.getChecksum())) { - throw new DataValidationException("Wrong otaPackage file!"); + if (!currentChecksum.equals(otaPackage.getChecksum())) { + throw new DataValidationException("Wrong otaPackage file!"); + } + } else { + //TODO: validate url } } @@ -264,6 +290,13 @@ public class BaseOtaPackageService implements OtaPackageService { if (otaPackageOld.getData() != null && !otaPackageOld.getData().equals(otaPackage.getData())) { throw new DataValidationException("Updating otaPackage data is prohibited!"); } + + if (otaPackageOld.getData() == null && otaPackage.getData() != null) { + DefaultTenantProfileConfiguration profileConfiguration = + (DefaultTenantProfileConfiguration) tenantProfileCache.get(tenantId).getProfileData().getConfiguration(); + long maxOtaPackagesInBytes = profileConfiguration.getMaxOtaPackagesInBytes(); + validateMaxSumDataSizePerTenant(tenantId, otaPackageDao, maxOtaPackagesInBytes, otaPackage.getDataSize(), OTA_PACKAGE); + } } }; diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java b/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java index 42f66663d1..ef8740030c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java @@ -16,8 +16,11 @@ package org.thingsboard.server.dao.ota; import org.thingsboard.server.common.data.OtaPackage; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.Dao; +import org.thingsboard.server.dao.TenantEntityDao; +import org.thingsboard.server.dao.TenantEntityWithDataDao; -public interface OtaPackageDao extends Dao { - +public interface OtaPackageDao extends Dao, TenantEntityWithDataDao { + Long sumDataSizeByTenantId(TenantId tenantId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageInfoDao.java b/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageInfoDao.java index d3294f0ec3..c40accf00a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageInfoDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageInfoDao.java @@ -28,7 +28,7 @@ public interface OtaPackageInfoDao extends Dao { PageData findOtaPackageInfoByTenantId(TenantId tenantId, PageLink pageLink); - PageData findOtaPackageInfoByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, boolean hasData, PageLink pageLink); + PageData findOtaPackageInfoByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, PageLink pageLink); boolean isOtaPackageUsed(OtaPackageId otaPackageId, OtaPackageType otaPackageType, DeviceProfileId deviceProfileId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java index 6dacd1357f..f442f5acea 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.hibernate.exception.ConstraintViolationException; +import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.ResourceType; import org.thingsboard.server.common.data.TbResource; @@ -28,16 +29,19 @@ import org.thingsboard.server.common.data.id.TbResourceId; 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.dao.exception.DataValidationException; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.Validator; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantDao; import java.util.List; import java.util.Optional; +import static org.thingsboard.server.common.data.EntityType.TB_RESOURCE; import static org.thingsboard.server.dao.device.DeviceServiceImpl.INCORRECT_TENANT_ID; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -49,12 +53,13 @@ public class BaseResourceService implements ResourceService { private final TbResourceDao resourceDao; private final TbResourceInfoDao resourceInfoDao; private final TenantDao tenantDao; + private final TbTenantProfileCache tenantProfileCache; - - public BaseResourceService(TbResourceDao resourceDao, TbResourceInfoDao resourceInfoDao, TenantDao tenantDao) { + public BaseResourceService(TbResourceDao resourceDao, TbResourceInfoDao resourceInfoDao, TenantDao tenantDao, @Lazy TbTenantProfileCache tenantProfileCache) { this.resourceDao = resourceDao; this.resourceInfoDao = resourceInfoDao; this.tenantDao = tenantDao; + this.tenantProfileCache = tenantProfileCache; } @Override @@ -143,8 +148,23 @@ public class BaseResourceService implements ResourceService { tenantResourcesRemover.removeEntities(tenantId, tenantId); } + @Override + public long sumDataSizeByTenantId(TenantId tenantId) { + return resourceDao.sumDataSizeByTenantId(tenantId); + } + private DataValidator resourceValidator = new DataValidator<>() { + @Override + protected void validateCreate(TenantId tenantId, TbResource resource) { + if (tenantId != null && !TenantId.SYS_TENANT_ID.equals(tenantId) ) { + DefaultTenantProfileConfiguration profileConfiguration = + (DefaultTenantProfileConfiguration) tenantProfileCache.get(tenantId).getProfileData().getConfiguration(); + long maxSumResourcesDataInBytes = profileConfiguration.getMaxResourcesInBytes(); + validateMaxSumDataSizePerTenant(tenantId, resourceDao, maxSumResourcesDataInBytes, resource.getData().length(), TB_RESOURCE); + } + } + @Override protected void validateDataImpl(TenantId tenantId, TbResource resource) { if (StringUtils.isEmpty(resource.getTitle())) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/TbResourceDao.java b/dao/src/main/java/org/thingsboard/server/dao/resource/TbResourceDao.java index 230e104191..f0f3f3a1e5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/resource/TbResourceDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/resource/TbResourceDao.java @@ -21,10 +21,11 @@ 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.dao.Dao; +import org.thingsboard.server.dao.TenantEntityWithDataDao; import java.util.List; -public interface TbResourceDao extends Dao { +public interface TbResourceDao extends Dao, TenantEntityWithDataDao { TbResource getResource(TenantId tenantId, ResourceType resourceType, String resourceId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java index ff8e79efc2..e630624ed4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java @@ -23,8 +23,10 @@ import org.hibernate.validator.cfg.ConstraintMapping; import org.thingsboard.server.common.data.BaseData; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.data.validation.NoXss; import org.thingsboard.server.dao.TenantEntityDao; +import org.thingsboard.server.dao.TenantEntityWithDataDao; import org.thingsboard.server.dao.exception.DataValidationException; import javax.validation.ConstraintViolation; @@ -123,6 +125,19 @@ public abstract class DataValidator> { } } + protected void validateMaxSumDataSizePerTenant(TenantId tenantId, + TenantEntityWithDataDao dataDao, + long maxSumDataSize, + long currentDataSize, + EntityType entityType) { + if (maxSumDataSize > 0) { + if (dataDao.sumDataSizeByTenantId(tenantId) + currentDataSize > maxSumDataSize) { + throw new DataValidationException(String.format("Failed to create the %s, files size limit is exhausted %d bytes!", + entityType.name().toLowerCase().replaceAll("_", " "), maxSumDataSize)); + } + } + } + protected static void validateJsonStructure(JsonNode expectedNode, JsonNode actualNode) { Set expectedFields = new HashSet<>(); Iterator fieldsIterator = expectedNode.fieldNames(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java index 6bad1af0e5..18852cbf11 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java @@ -26,6 +26,7 @@ import org.thingsboard.server.dao.Dao; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.model.BaseEntity; +import java.util.Collection; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -87,6 +88,12 @@ public abstract class JpaAbstractDao, D> return !getCrudRepository().existsById(id); } + @Transactional + public void removeAllByIds(Collection ids) { + CrudRepository repository = getCrudRepository(); + ids.forEach(repository::deleteById); + } + @Override public List find(TenantId tenantId) { List entities = Lists.newArrayList(getCrudRepository().findAll()); 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 b4c0ac09c6..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 @@ -17,11 +17,13 @@ package org.thingsboard.server.dao.sql.alarm; import org.springframework.data.domain.Page; import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.Modifying; 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.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; @@ -159,4 +161,8 @@ public interface AlarmRepository extends CrudRepository { @Param("affectedEntityId") UUID affectedEntityId, @Param("affectedEntityType") String affectedEntityType, @Param("alarmStatuses") Set alarmStatuses); + + @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 cf222413b0..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,10 +26,12 @@ 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; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.dao.DaoUtil; @@ -161,4 +163,10 @@ public class JpaAlarmDao extends JpaAbstractDao implements A public Set findAlarmSeverities(TenantId tenantId, EntityId entityId, Set statuses) { return alarmRepository.findAlarmSeverities(tenantId.getId(), entityId.getId(), entityId.getEntityType().name(), statuses); } + + @Override + 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/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java index 95737ca48d..98309b9e51 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java @@ -20,6 +20,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.repository.CrudRepository; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.OtaPackage; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.ota.OtaPackageDao; import org.thingsboard.server.dao.model.sql.OtaPackageEntity; import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao; @@ -43,4 +44,8 @@ public class JpaOtaPackageDao extends JpaAbstractSearchTextDao findOtaPackageInfoByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, boolean hasData, PageLink pageLink) { + public PageData findOtaPackageInfoByTenantIdAndDeviceProfileIdAndTypeAndHasData(TenantId tenantId, DeviceProfileId deviceProfileId, OtaPackageType otaPackageType, PageLink pageLink) { return DaoUtil.toPageData(otaPackageInfoRepository .findAllByTenantIdAndTypeAndDeviceProfileIdAndHasData( tenantId.getId(), deviceProfileId.getId(), otaPackageType, - hasData, Objects.toString(pageLink.getTextSearch(), ""), DaoUtil.toPageable(pageLink))); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageInfoRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageInfoRepository.java index 9848f83200..b380f8a150 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageInfoRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageInfoRepository.java @@ -26,27 +26,26 @@ import org.thingsboard.server.dao.model.sql.OtaPackageInfoEntity; import java.util.UUID; public interface OtaPackageInfoRepository extends CrudRepository { - @Query("SELECT new OtaPackageInfoEntity(f.id, f.createdTime, f.tenantId, f.deviceProfileId, f.type, f.title, f.version, f.fileName, f.contentType, f.checksumAlgorithm, f.checksum, f.dataSize, f.additionalInfo, f.data IS NOT NULL) FROM OtaPackageEntity f WHERE " + + @Query("SELECT new OtaPackageInfoEntity(f.id, f.createdTime, f.tenantId, f.deviceProfileId, f.type, f.title, f.version, f.url, f.fileName, f.contentType, f.checksumAlgorithm, f.checksum, f.dataSize, f.additionalInfo, CASE WHEN (f.data IS NOT NULL OR f.url IS NOT NULL) THEN true ELSE false END) FROM OtaPackageEntity f WHERE " + "f.tenantId = :tenantId " + "AND LOWER(f.searchText) LIKE LOWER(CONCAT(:searchText, '%'))") Page findAllByTenantId(@Param("tenantId") UUID tenantId, @Param("searchText") String searchText, Pageable pageable); - @Query("SELECT new OtaPackageInfoEntity(f.id, f.createdTime, f.tenantId, f.deviceProfileId, f.type, f.title, f.version, f.fileName, f.contentType, f.checksumAlgorithm, f.checksum, f.dataSize, f.additionalInfo, f.data IS NOT NULL) FROM OtaPackageEntity f WHERE " + + @Query("SELECT new OtaPackageInfoEntity(f.id, f.createdTime, f.tenantId, f.deviceProfileId, f.type, f.title, f.version, f.url, f.fileName, f.contentType, f.checksumAlgorithm, f.checksum, f.dataSize, f.additionalInfo, true) FROM OtaPackageEntity f WHERE " + "f.tenantId = :tenantId " + "AND f.deviceProfileId = :deviceProfileId " + "AND f.type = :type " + - "AND ((f.data IS NOT NULL AND :hasData = true) OR (f.data IS NULL AND :hasData = false ))" + + "AND (f.data IS NOT NULL OR f.url IS NOT NULL) " + "AND LOWER(f.searchText) LIKE LOWER(CONCAT(:searchText, '%'))") Page findAllByTenantIdAndTypeAndDeviceProfileIdAndHasData(@Param("tenantId") UUID tenantId, @Param("deviceProfileId") UUID deviceProfileId, @Param("type") OtaPackageType type, - @Param("hasData") boolean hasData, @Param("searchText") String searchText, Pageable pageable); - @Query("SELECT new OtaPackageInfoEntity(f.id, f.createdTime, f.tenantId, f.deviceProfileId, f.type, f.title, f.version, f.fileName, f.contentType, f.checksumAlgorithm, f.checksum, f.dataSize, f.additionalInfo, f.data IS NOT NULL) FROM OtaPackageEntity f WHERE f.id = :id") + @Query("SELECT new OtaPackageInfoEntity(f.id, f.createdTime, f.tenantId, f.deviceProfileId, f.type, f.title, f.version, f.url, f.fileName, f.contentType, f.checksumAlgorithm, f.checksum, f.dataSize, f.additionalInfo, CASE WHEN (f.data IS NOT NULL OR f.url IS NOT NULL) THEN true ELSE false END) FROM OtaPackageEntity f WHERE f.id = :id") OtaPackageInfoEntity findOtaPackageInfoById(@Param("id") UUID id); @Query(value = "SELECT exists(SELECT * " + diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageRepository.java index 3699005ff2..47e6f82285 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/OtaPackageRepository.java @@ -15,10 +15,15 @@ */ package org.thingsboard.server.dao.sql.ota; +import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.CrudRepository; +import org.springframework.data.repository.query.Param; import org.thingsboard.server.dao.model.sql.OtaPackageEntity; +import org.thingsboard.server.dao.model.sql.OtaPackageInfoEntity; import java.util.UUID; public interface OtaPackageRepository extends CrudRepository { + @Query(value = "SELECT COALESCE(SUM(ota.data_size), 0) FROM ota_package ota WHERE ota.tenant_id = :tenantId AND ota.data IS NOT NULL", nativeQuery = true) + Long sumDataSizeByTenantId(@Param("tenantId") UUID tenantId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/resource/JpaTbResourceDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/resource/JpaTbResourceDao.java index f35f654f77..324549765e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/resource/JpaTbResourceDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/resource/JpaTbResourceDao.java @@ -92,4 +92,8 @@ public class JpaTbResourceDao extends JpaAbstractSearchTextDao Objects.toString(pageLink.getTextSearch(), ""), DaoUtil.toPageable(pageLink, TenantInfoEntity.tenantInfoColumnMap))); } + + @Override + public PageData findTenantsIds(PageLink pageLink) { + return DaoUtil.pageToPageData(tenantRepository.findTenantsIds(DaoUtil.toPageable(pageLink))).mapData(TenantId::new); + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java index b43d70197c..8ab12e0bb5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java @@ -50,4 +50,8 @@ public interface TenantRepository extends PagingAndSortingRepository findTenantInfoByRegionNextPage(@Param("region") String region, @Param("textSearch") String textSearch, Pageable pageable); + + @Query("SELECT t.id FROM TenantEntity t") + Page findTenantsIds(Pageable pageable); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java index bff1c2c8a9..5bceb35376 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java @@ -46,5 +46,7 @@ public interface TenantDao extends Dao { PageData findTenantsByRegion(TenantId tenantId, String region, PageLink pageLink); PageData findTenantInfosByRegion(TenantId tenantId, String region, PageLink pageLink); - + + PageData findTenantsIds(PageLink pageLink); + } diff --git a/dao/src/main/resources/sql/schema-entities-hsql.sql b/dao/src/main/resources/sql/schema-entities-hsql.sql index dfc1821f66..2c5113867a 100644 --- a/dao/src/main/resources/sql/schema-entities-hsql.sql +++ b/dao/src/main/resources/sql/schema-entities-hsql.sql @@ -168,6 +168,7 @@ CREATE TABLE IF NOT EXISTS ota_package ( type varchar(32) NOT NULL, title varchar(255) NOT NULL, version varchar(255) NOT NULL, + url varchar(255), file_name varchar(255), content_type varchar(255), checksum_algorithm varchar(32), diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index be7e836a65..3140ff83d6 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -186,6 +186,7 @@ CREATE TABLE IF NOT EXISTS ota_package ( type varchar(32) NOT NULL, title varchar(255) NOT NULL, version varchar(255) NOT NULL, + url varchar(255), file_name varchar(255), content_type varchar(255), checksum_algorithm varchar(32), diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java index 2aeb44c0be..022051fe8c 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java @@ -59,6 +59,7 @@ import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.settings.AdminSettingsService; +import org.thingsboard.server.dao.tenant.DefaultTbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantProfileService; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -158,10 +159,12 @@ public abstract class AbstractServiceTest { @Autowired protected ResourceService resourceService; - @Autowired protected OtaPackageService otaPackageService; + @Autowired + protected DefaultTbTenantProfileCache tenantProfileCache; + public class IdComparator implements Comparator { @Override public int compare(D o1, D o2) { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseOtaPackageServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseOtaPackageServiceTest.java index ab895e1052..35063eef75 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseOtaPackageServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseOtaPackageServiceTest.java @@ -28,11 +28,13 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.Tenant; -import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; 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.dao.exception.DataValidationException; import java.nio.ByteBuffer; @@ -50,7 +52,9 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { private static final String CONTENT_TYPE = "text/plain"; private static final ChecksumAlgorithm CHECKSUM_ALGORITHM = ChecksumAlgorithm.SHA256; private static final String CHECKSUM = "4bf5122f344554c53bde2ebb8cd2b7e3d1600ad631c385a5d7cce23c7785459a"; - private static final ByteBuffer DATA = ByteBuffer.wrap(new byte[]{1}); + private static final long DATA_SIZE = 1L; + private static final ByteBuffer DATA = ByteBuffer.wrap(new byte[]{(int) DATA_SIZE}); + private static final String URL = "http://firmware.test.org"; private IdComparator idComparator = new IdComparator<>(); @@ -78,6 +82,41 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { @After public void after() { tenantService.deleteTenant(tenantId); + tenantProfileService.deleteTenantProfiles(tenantId); + } + + @Test + public void testSaveOtaPackageWithMaxSumDataSizeOutOfLimit() { + TenantProfile defaultTenantProfile = tenantProfileService.findDefaultTenantProfile(tenantId); + defaultTenantProfile.getProfileData().setConfiguration(DefaultTenantProfileConfiguration.builder().maxOtaPackagesInBytes(DATA_SIZE).build()); + tenantProfileService.saveTenantProfile(tenantId, defaultTenantProfile); + + Assert.assertEquals(0, otaPackageService.sumDataSizeByTenantId(tenantId)); + + createFirmware(tenantId, "1"); + Assert.assertEquals(1, otaPackageService.sumDataSizeByTenantId(tenantId)); + + thrown.expect(DataValidationException.class); + thrown.expectMessage(String.format("Failed to create the ota package, files size limit is exhausted %d bytes!", DATA_SIZE)); + createFirmware(tenantId, "2"); + } + + @Test + public void sumDataSizeByTenantId() { + Assert.assertEquals(0, otaPackageService.sumDataSizeByTenantId(tenantId)); + + createFirmware(tenantId, "0.1"); + Assert.assertEquals(1, otaPackageService.sumDataSizeByTenantId(tenantId)); + + int maxSumDataSize = 8; + List packages = new ArrayList<>(maxSumDataSize); + + for (int i = 2; i <= maxSumDataSize; i++) { + packages.add(createFirmware(tenantId, "0." + i)); + Assert.assertEquals(i, otaPackageService.sumDataSizeByTenantId(tenantId)); + } + + Assert.assertEquals(maxSumDataSize, otaPackageService.sumDataSizeByTenantId(tenantId)); } @Test @@ -93,6 +132,7 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); firmware.setChecksum(CHECKSUM); firmware.setData(DATA); + firmware.setDataSize(DATA_SIZE); OtaPackage savedFirmware = otaPackageService.saveOtaPackage(firmware); Assert.assertNotNull(savedFirmware); @@ -113,6 +153,35 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { otaPackageService.deleteOtaPackage(tenantId, savedFirmware.getId()); } + @Test + public void testSaveFirmwareWithUrl() { + OtaPackageInfo firmware = new OtaPackageInfo(); + firmware.setTenantId(tenantId); + firmware.setDeviceProfileId(deviceProfileId); + firmware.setType(FIRMWARE); + firmware.setTitle(TITLE); + firmware.setVersion(VERSION); + firmware.setUrl(URL); + firmware.setDataSize(0L); + OtaPackageInfo savedFirmware = otaPackageService.saveOtaPackageInfo(firmware); + + Assert.assertNotNull(savedFirmware); + Assert.assertNotNull(savedFirmware.getId()); + Assert.assertTrue(savedFirmware.getCreatedTime() > 0); + Assert.assertEquals(firmware.getTenantId(), savedFirmware.getTenantId()); + Assert.assertEquals(firmware.getTitle(), savedFirmware.getTitle()); + Assert.assertEquals(firmware.getFileName(), savedFirmware.getFileName()); + Assert.assertEquals(firmware.getContentType(), savedFirmware.getContentType()); + + savedFirmware.setAdditionalInfo(JacksonUtil.newObjectNode()); + otaPackageService.saveOtaPackageInfo(savedFirmware); + + OtaPackage foundFirmware = otaPackageService.findOtaPackageById(tenantId, savedFirmware.getId()); + Assert.assertEquals(foundFirmware.getTitle(), savedFirmware.getTitle()); + + otaPackageService.deleteOtaPackage(tenantId, savedFirmware.getId()); + } + @Test public void testSaveFirmwareInfoAndUpdateWithData() { OtaPackageInfo firmwareInfo = new OtaPackageInfo(); @@ -141,6 +210,7 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); firmware.setChecksum(CHECKSUM); firmware.setData(DATA); + firmware.setDataSize(DATA_SIZE); otaPackageService.saveOtaPackage(firmware); @@ -345,50 +415,15 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { @Test public void testSaveFirmwareWithExistingTitleAndVersion() { - OtaPackage firmware = new OtaPackage(); - firmware.setTenantId(tenantId); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(TITLE); - firmware.setVersion(VERSION); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - otaPackageService.saveOtaPackage(firmware); - - OtaPackage newFirmware = new OtaPackage(); - newFirmware.setTenantId(tenantId); - newFirmware.setDeviceProfileId(deviceProfileId); - newFirmware.setType(FIRMWARE); - newFirmware.setTitle(TITLE); - newFirmware.setVersion(VERSION); - newFirmware.setFileName(FILE_NAME); - newFirmware.setContentType(CONTENT_TYPE); - newFirmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - newFirmware.setChecksum(CHECKSUM); - newFirmware.setData(DATA); - + createFirmware(tenantId, VERSION); thrown.expect(DataValidationException.class); thrown.expectMessage("OtaPackage with such title and version already exists!"); - otaPackageService.saveOtaPackage(newFirmware); + createFirmware(tenantId, VERSION); } @Test public void testDeleteFirmwareWithReferenceByDevice() { - OtaPackage firmware = new OtaPackage(); - firmware.setTenantId(tenantId); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(TITLE); - firmware.setVersion(VERSION); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - OtaPackage savedFirmware = otaPackageService.saveOtaPackage(firmware); + OtaPackage savedFirmware = createFirmware(tenantId, VERSION); Device device = new Device(); device.setTenantId(tenantId); @@ -409,18 +444,7 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { @Test public void testUpdateDeviceProfileId() { - OtaPackage firmware = new OtaPackage(); - firmware.setTenantId(tenantId); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(TITLE); - firmware.setVersion(VERSION); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - OtaPackage savedFirmware = otaPackageService.saveOtaPackage(firmware); + OtaPackage savedFirmware = createFirmware(tenantId, VERSION); try { thrown.expect(DataValidationException.class); @@ -448,6 +472,7 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); firmware.setChecksum(CHECKSUM); firmware.setData(DATA); + firmware.setDataSize(DATA_SIZE); OtaPackage savedFirmware = otaPackageService.saveOtaPackage(firmware); savedDeviceProfile.setFirmwareId(savedFirmware.getId()); @@ -465,18 +490,7 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { @Test public void testFindFirmwareById() { - OtaPackage firmware = new OtaPackage(); - firmware.setTenantId(tenantId); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(TITLE); - firmware.setVersion(VERSION); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - OtaPackage savedFirmware = otaPackageService.saveOtaPackage(firmware); + OtaPackage savedFirmware = createFirmware(tenantId, VERSION); OtaPackage foundFirmware = otaPackageService.findOtaPackageById(tenantId, savedFirmware.getId()); Assert.assertNotNull(foundFirmware); @@ -502,18 +516,7 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { @Test public void testDeleteFirmware() { - OtaPackage firmware = new OtaPackage(); - firmware.setTenantId(tenantId); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(TITLE); - firmware.setVersion(VERSION); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - OtaPackage savedFirmware = otaPackageService.saveOtaPackage(firmware); + OtaPackage savedFirmware = createFirmware(tenantId, VERSION); OtaPackage foundFirmware = otaPackageService.findOtaPackageById(tenantId, savedFirmware.getId()); Assert.assertNotNull(foundFirmware); @@ -526,23 +529,25 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { public void testFindTenantFirmwaresByTenantId() { List firmwares = new ArrayList<>(); for (int i = 0; i < 165; i++) { - OtaPackage firmware = new OtaPackage(); - firmware.setTenantId(tenantId); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(TITLE); - firmware.setVersion(VERSION + i); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - - OtaPackageInfo info = new OtaPackageInfo(otaPackageService.saveOtaPackage(firmware)); + OtaPackageInfo info = new OtaPackageInfo(createFirmware(tenantId, VERSION + i)); info.setHasData(true); firmwares.add(info); } + OtaPackageInfo firmwareWithUrl = new OtaPackageInfo(); + firmwareWithUrl.setTenantId(tenantId); + firmwareWithUrl.setDeviceProfileId(deviceProfileId); + firmwareWithUrl.setType(FIRMWARE); + firmwareWithUrl.setTitle(TITLE); + firmwareWithUrl.setVersion(VERSION); + firmwareWithUrl.setUrl(URL); + firmwareWithUrl.setDataSize(0L); + + OtaPackageInfo savedFwWithUrl = otaPackageService.saveOtaPackageInfo(firmwareWithUrl); + savedFwWithUrl.setHasData(true); + + firmwares.add(savedFwWithUrl); + List loadedFirmwares = new ArrayList<>(); PageLink pageLink = new PageLink(16); PageData pageData; @@ -571,58 +576,38 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { public void testFindTenantFirmwaresByTenantIdAndHasData() { List firmwares = new ArrayList<>(); for (int i = 0; i < 165; i++) { - OtaPackageInfo firmwareInfo = new OtaPackageInfo(); - firmwareInfo.setTenantId(tenantId); - firmwareInfo.setDeviceProfileId(deviceProfileId); - firmwareInfo.setType(FIRMWARE); - firmwareInfo.setTitle(TITLE); - firmwareInfo.setVersion(VERSION + i); - firmwareInfo.setFileName(FILE_NAME); - firmwareInfo.setContentType(CONTENT_TYPE); - firmwareInfo.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmwareInfo.setChecksum(CHECKSUM); - firmwareInfo.setDataSize((long) DATA.array().length); - firmwares.add(otaPackageService.saveOtaPackageInfo(firmwareInfo)); + firmwares.add(new OtaPackageInfo(otaPackageService.saveOtaPackage(createFirmware(tenantId, VERSION + i)))); } + OtaPackageInfo firmwareWithUrl = new OtaPackageInfo(); + firmwareWithUrl.setTenantId(tenantId); + firmwareWithUrl.setDeviceProfileId(deviceProfileId); + firmwareWithUrl.setType(FIRMWARE); + firmwareWithUrl.setTitle(TITLE); + firmwareWithUrl.setVersion(VERSION); + firmwareWithUrl.setUrl(URL); + firmwareWithUrl.setDataSize(0L); + + OtaPackageInfo savedFwWithUrl = otaPackageService.saveOtaPackageInfo(firmwareWithUrl); + savedFwWithUrl.setHasData(true); + + firmwares.add(savedFwWithUrl); + List loadedFirmwares = new ArrayList<>(); PageLink pageLink = new PageLink(16); PageData pageData; do { - pageData = otaPackageService.findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(tenantId, deviceProfileId, FIRMWARE, false, pageLink); + pageData = otaPackageService.findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(tenantId, deviceProfileId, FIRMWARE, pageLink); loadedFirmwares.addAll(pageData.getData()); if (pageData.hasNext()) { pageLink = pageLink.nextPageLink(); } } while (pageData.hasNext()); - Collections.sort(firmwares, idComparator); - Collections.sort(loadedFirmwares, idComparator); - - Assert.assertEquals(firmwares, loadedFirmwares); - - firmwares.forEach(f -> { - OtaPackage firmware = new OtaPackage(f.getId()); - firmware.setCreatedTime(f.getCreatedTime()); - firmware.setTenantId(f.getTenantId()); - firmware.setDeviceProfileId(deviceProfileId); - firmware.setType(FIRMWARE); - firmware.setTitle(f.getTitle()); - firmware.setVersion(f.getVersion()); - firmware.setFileName(FILE_NAME); - firmware.setContentType(CONTENT_TYPE); - firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); - firmware.setChecksum(CHECKSUM); - firmware.setData(DATA); - firmware.setDataSize((long) DATA.array().length); - otaPackageService.saveOtaPackage(firmware); - f.setHasData(true); - }); - loadedFirmwares = new ArrayList<>(); pageLink = new PageLink(16); do { - pageData = otaPackageService.findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(tenantId, deviceProfileId, FIRMWARE, true, pageLink); + pageData = otaPackageService.findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(tenantId, deviceProfileId, FIRMWARE, pageLink); loadedFirmwares.addAll(pageData.getData()); if (pageData.hasNext()) { pageLink = pageLink.nextPageLink(); @@ -642,4 +627,20 @@ public abstract class BaseOtaPackageServiceTest extends AbstractServiceTest { Assert.assertTrue(pageData.getData().isEmpty()); } + private OtaPackage createFirmware(TenantId tenantId, String version) { + OtaPackage firmware = new OtaPackage(); + firmware.setTenantId(tenantId); + firmware.setDeviceProfileId(deviceProfileId); + firmware.setType(FIRMWARE); + firmware.setTitle(TITLE); + firmware.setVersion(version); + firmware.setFileName(FILE_NAME); + firmware.setContentType(CONTENT_TYPE); + firmware.setChecksumAlgorithm(CHECKSUM_ALGORITHM); + firmware.setChecksum(CHECKSUM); + firmware.setData(DATA); + firmware.setDataSize(DATA_SIZE); + return otaPackageService.saveOtaPackage(firmware); + } + } diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 2f4fd8751d..bec9a47b44 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -24,7 +24,7 @@ const topicProperties = config.get('kafka.topic_properties'); const kafkaClientId = config.get('kafka.client_id'); const acks = Number(config.get('kafka.acks')); const requestTimeout = Number(config.get('kafka.requestTimeout')); -const compressionType = (config.get('kafka.requestTimeout') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; +const compressionType = (config.get('kafka.compression') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; let kafkaClient; let kafkaAdmin; diff --git a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java index 433e2f8d96..8ad95805a5 100644 --- a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java +++ b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java @@ -19,18 +19,25 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.io.ByteArrayResource; +import org.springframework.core.io.Resource; import org.springframework.http.HttpEntity; +import org.springframework.http.HttpHeaders; import org.springframework.http.HttpMethod; import org.springframework.http.HttpRequest; import org.springframework.http.HttpStatus; +import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.http.client.ClientHttpRequestExecution; import org.springframework.http.client.ClientHttpRequestInterceptor; import org.springframework.http.client.ClientHttpResponse; import org.springframework.http.client.support.HttpRequestWrapper; +import org.springframework.util.LinkedMultiValueMap; +import org.springframework.util.MultiValueMap; import org.springframework.util.StringUtils; import org.springframework.web.client.HttpClientErrorException; import org.springframework.web.client.RestTemplate; +import org.springframework.web.multipart.MultipartFile; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.rest.client.utils.RestJsonConverter; import org.thingsboard.server.common.data.AdminSettings; @@ -48,6 +55,10 @@ import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityViewInfo; import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.OtaPackage; +import org.thingsboard.server.common.data.OtaPackageInfo; +import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantInfo; import org.thingsboard.server.common.data.TenantProfile; @@ -78,8 +89,10 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.OAuth2ClientRegistrationTemplateId; +import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; +import org.thingsboard.server.common.data.id.TbResourceId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.id.UserId; @@ -91,6 +104,8 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.oauth2.OAuth2ClientInfo; import org.thingsboard.server.common.data.oauth2.OAuth2ClientRegistrationTemplate; import org.thingsboard.server.common.data.oauth2.OAuth2ClientsParams; +import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; +import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.SortOrder; @@ -127,9 +142,9 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.stream.Collectors; @@ -147,7 +162,6 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { private final ObjectMapper objectMapper = new ObjectMapper(); private ExecutorService service = ThingsBoardExecutors.newWorkStealingPool(10, getClass()); - protected static final String ACTIVATE_TOKEN_REGEX = "/api/noauth/activate?activateToken="; public RestClient(String baseURL) { @@ -1238,6 +1252,21 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { HttpEntity.EMPTY, Device.class, tenantId, deviceId).getBody(); } + public Long countDevicesByTenantIdAndDeviceProfileIdAndEmptyOtaPackage(OtaPackageType otaPackageType, DeviceProfileId deviceProfileId) { + Map params = new HashMap<>(); + params.put("otaPackageType", otaPackageType.name()); + params.put("deviceProfileId", deviceProfileId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/devices/count/{otaPackageType}?deviceProfileId={deviceProfileId}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference() { + }, + params + ).getBody(); + } + @Deprecated public Device createDevice(String name, String type) { Device device = new Device(); @@ -2830,6 +2859,176 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { restTemplate.postForEntity(baseURL + "/api/edge/sync/{edgeId}", null, EdgeId.class, params); } + public ResponseEntity downloadResource(TbResourceId resourceId) { + Map params = new HashMap<>(); + params.put("resourceId", resourceId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/resource/{resourceId}/download", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference<>() {}, + params + ); + } + + public TbResourceInfo getResourceInfoById(TbResourceId resourceId) { + Map params = new HashMap<>(); + params.put("resourceId", resourceId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/resource/info/{resourceId}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference() {}, + params + ).getBody(); + } + + public TbResource getResourceId(TbResourceId resourceId) { + Map params = new HashMap<>(); + params.put("resourceId", resourceId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/resource/{resourceId}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference() {}, + params + ).getBody(); + } + + public TbResource saveResource(TbResource resource) { + return restTemplate.postForEntity( + baseURL + "/api/resource", + resource, + TbResource.class + ).getBody(); + } + + public PageData getResources(PageLink pageLink) { + Map params = new HashMap<>(); + addPageLinkToParam(params, pageLink); + return restTemplate.exchange( + baseURL + "/api/resource?" + getUrlParams(pageLink), + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference>() {}, + params + ).getBody(); + } + + public void deleteResource(TbResourceId resourceId) { + restTemplate.delete("/api/resource/{resourceId}", resourceId.getId().toString()); + } + + public ResponseEntity downloadOtaPackage(OtaPackageId otaPackageId) { + Map params = new HashMap<>(); + params.put("otaPackageId", otaPackageId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/otaPackage/{otaPackageId}/download", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference<>() {}, + params + ); + } + + public OtaPackageInfo getOtaPackageInfoById(OtaPackageId otaPackageId) { + Map params = new HashMap<>(); + params.put("otaPackageId", otaPackageId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/otaPackage/info/{otaPackageId}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference() {}, + params + ).getBody(); + } + + public OtaPackage getOtaPackageById(OtaPackageId otaPackageId) { + Map params = new HashMap<>(); + params.put("otaPackageId", otaPackageId.getId().toString()); + + return restTemplate.exchange( + baseURL + "/api/otaPackage/{otaPackageId}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference() {}, + params + ).getBody(); + } + + public OtaPackageInfo saveOtaPackageInfo(OtaPackageInfo otaPackageInfo) { + return restTemplate.postForEntity(baseURL + "/api/otaPackage", otaPackageInfo, OtaPackageInfo.class).getBody(); + } + + public OtaPackage saveOtaPackageData(OtaPackageId otaPackageId, String checkSum, ChecksumAlgorithm checksumAlgorithm, MultipartFile file) throws Exception { + HttpHeaders header = new HttpHeaders(); + header.setContentType(MediaType.MULTIPART_FORM_DATA); + + MultiValueMap fileMap = new LinkedMultiValueMap<>(); + fileMap.add(HttpHeaders.CONTENT_DISPOSITION, "form-data; name=file; filename=" + file.getName()); + HttpEntity fileEntity = new HttpEntity<>(new ByteArrayResource(file.getBytes()), fileMap); + + MultiValueMap body = new LinkedMultiValueMap<>(); + body.add("file", fileEntity); + HttpEntity> requestEntity = new HttpEntity<>(body, header); + + Map params = new HashMap<>(); + params.put("otaPackageId", otaPackageId.getId().toString()); + params.put("checksumAlgorithm", checksumAlgorithm.name()); + String url = "/api/otaPackage/{otaPackageId}?checksumAlgorithm={checksumAlgorithm}"; + + if(checkSum != null) { + url += "&checkSum={checkSum}"; + } + + return restTemplate.postForEntity( + baseURL + url, requestEntity, OtaPackage.class, params + ).getBody(); + } + + public PageData getOtaPackages(PageLink pageLink) { + Map params = new HashMap<>(); + addPageLinkToParam(params, pageLink); + + return restTemplate.exchange( + baseURL + "/api/otaPackages?" + getUrlParams(pageLink), + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference>() { + }, + params + ).getBody(); + } + + public PageData getOtaPackages(DeviceProfileId deviceProfileId, + OtaPackageType otaPackageType, + boolean hasData, + PageLink pageLink) { + Map params = new HashMap<>(); + params.put("hasData", String.valueOf(hasData)); + params.put("deviceProfileId", deviceProfileId.getId().toString()); + params.put("type", otaPackageType.name()); + addPageLinkToParam(params, pageLink); + + return restTemplate.exchange( + baseURL + "/api/otaPackages/{deviceProfileId}/{type}/{hasData}?" + getUrlParams(pageLink), + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference>() { + }, + params + ).getBody(); + } + + public void deleteOtaPackage(OtaPackageId otaPackageId) { + restTemplate.delete(baseURL + "/api/otaPackage/{otaPackageId}", otaPackageId.getId().toString()); + } + @Deprecated public Optional getAttributes(String accessToken, String clientKeys, String sharedKeys) { Map params = new HashMap<>(); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java index 7a0e3cce63..43e3a5e329 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java @@ -47,7 +47,9 @@ import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.nosql.CassandraStatementTask; import org.thingsboard.server.dao.nosql.TbResultSetFuture; +import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -202,6 +204,10 @@ public interface TbContext { EntityViewService getEntityViewService(); + ResourceService getResourceService(); + + OtaPackageService getOtaPackageService(); + RuleEngineDeviceProfileCache getDeviceProfileCache(); EdgeService getEdgeService(); 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..6b7a695eea 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 @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.device.profile.AlarmConditionFilterKey; import org.thingsboard.server.common.data.device.profile.AlarmConditionKeyType; import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm; +import org.thingsboard.server.common.data.exception.ApiUsageLimitsExceededException; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; @@ -150,6 +151,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 +196,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)) { @@ -253,7 +262,12 @@ class DeviceState { for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { AlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), a -> new AlarmState(this.deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm), dynamicPredicateValueCtx)); - stateChanged |= alarmState.process(ctx, msg, latestValues, update); + try { + stateChanged |= alarmState.process(ctx, msg, latestValues, update); + } catch (ApiUsageLimitsExceededException e) { + alarmStates.remove(alarm.getId()); + throw e; + } } } } diff --git a/ui-ngx/src/app/core/api/widget-api.models.ts b/ui-ngx/src/app/core/api/widget-api.models.ts index 8c6ac47ba5..515e5397f1 100644 --- a/ui-ngx/src/app/core/api/widget-api.models.ts +++ b/ui-ngx/src/app/core/api/widget-api.models.ts @@ -152,6 +152,7 @@ export interface IStateController { getStateIndex(): number; getStateIdAtIndex(index: number): string; getEntityId(entityParamName: string): EntityId; + getCurrentStateName(): string; } export interface SubscriptionInfo { diff --git a/ui-ngx/src/app/core/auth/auth.service.ts b/ui-ngx/src/app/core/auth/auth.service.ts index 4d25e5ce72..c17d993bf7 100644 --- a/ui-ngx/src/app/core/auth/auth.service.ts +++ b/ui-ngx/src/app/core/auth/auth.service.ts @@ -45,6 +45,7 @@ import { ActionNotificationShow } from '@core/notification/notification.actions' import { MatDialog, MatDialogConfig } from '@angular/material/dialog'; import { AlertDialogComponent } from '@shared/components/dialog/alert-dialog.component'; import { OAuth2ClientInfo } from '@shared/models/oauth2.models'; +import { isMobileApp } from '@core/utils'; @Injectable({ providedIn: 'root' @@ -194,11 +195,13 @@ export class AuthService { } public gotoDefaultPlace(isAuthenticated: boolean) { - const authState = getCurrentAuthState(this.store); - const url = this.defaultUrl(isAuthenticated, authState); - this.zone.run(() => { - this.router.navigateByUrl(url); - }); + if (!isMobileApp()) { + const authState = getCurrentAuthState(this.store); + const url = this.defaultUrl(isAuthenticated, authState); + this.zone.run(() => { + this.router.navigateByUrl(url); + }); + } } public loadOAuth2Clients(): Observable> { @@ -516,12 +519,15 @@ export class AuthService { return this.refreshTokenSubject !== null; } - public setUserFromJwtToken(jwtToken, refreshToken, notify) { + public setUserFromJwtToken(jwtToken, refreshToken, notify): Observable { + const authenticatedSubject = new ReplaySubject(); if (!jwtToken) { AuthService.clearTokenData(); if (notify) { this.notifyUnauthenticated(); } + authenticatedSubject.next(false); + authenticatedSubject.complete(); } else { this.updateAndValidateTokens(jwtToken, refreshToken, true); if (notify) { @@ -530,16 +536,30 @@ export class AuthService { (authPayload) => { this.notifyUserLoaded(true); this.notifyAuthenticated(authPayload); + authenticatedSubject.next(true); + authenticatedSubject.complete(); }, () => { this.notifyUserLoaded(true); this.notifyUnauthenticated(); + authenticatedSubject.next(false); + authenticatedSubject.complete(); } ); } else { - this.loadUser(false).subscribe(); + this.loadUser(false).subscribe( + () => { + authenticatedSubject.next(true); + authenticatedSubject.complete(); + }, + () => { + authenticatedSubject.next(false); + authenticatedSubject.complete(); + } + ); } } + return authenticatedSubject; } private updateAndValidateTokens(jwtToken, refreshToken, notify: boolean) { diff --git a/ui-ngx/src/app/core/guards/auth.guard.ts b/ui-ngx/src/app/core/guards/auth.guard.ts index fdba07b94a..98fd4bd969 100644 --- a/ui-ngx/src/app/core/guards/auth.guard.ts +++ b/ui-ngx/src/app/core/guards/auth.guard.ts @@ -29,6 +29,7 @@ import { DialogService } from '@core/services/dialog.service'; import { TranslateService } from '@ngx-translate/core'; import { UtilsService } from '@core/services/utils.service'; import { isObject } from '@core/utils'; +import { MobileService } from '@core/services/mobile.service'; @Injectable({ providedIn: 'root' @@ -41,6 +42,7 @@ export class AuthGuard implements CanActivate, CanActivateChild { private dialogService: DialogService, private utils: UtilsService, private translate: TranslateService, + private mobileService: MobileService, private zone: NgZone) {} getAuthState(): Observable { @@ -108,6 +110,10 @@ export class AuthGuard implements CanActivate, CanActivateChild { return of(false); } } + if (this.mobileService.isMobileApp() && !path.startsWith('dashboard.')) { + this.mobileService.handleMobileNavigation(path, params); + return of(false); + } const defaultUrl = this.authService.defaultUrl(true, authState, path, params); if (defaultUrl) { // this.authService.gotoDefaultPlace(true); diff --git a/ui-ngx/src/app/core/http/ota-package.service.ts b/ui-ngx/src/app/core/http/ota-package.service.ts index 3ae1490a88..0042f67b24 100644 --- a/ui-ngx/src/app/core/http/ota-package.service.ts +++ b/ui-ngx/src/app/core/http/ota-package.service.ts @@ -40,7 +40,7 @@ export class OtaPackageService { public getOtaPackagesInfoByDeviceProfileId(pageLink: PageLink, deviceProfileId: string, type: OtaUpdateType, hasData = true, config?: RequestConfig): Observable> { - const url = `/api/otaPackages/${deviceProfileId}/${type}/${hasData}${pageLink.toQuery()}`; + const url = `/api/otaPackages/${deviceProfileId}/${type}${pageLink.toQuery()}`; return this.http.get>(url, defaultHttpOptionsFromConfig(config)); } diff --git a/ui-ngx/src/app/core/services/mobile.service.ts b/ui-ngx/src/app/core/services/mobile.service.ts index b6774f37d0..800651356a 100644 --- a/ui-ngx/src/app/core/services/mobile.service.ts +++ b/ui-ngx/src/app/core/services/mobile.service.ts @@ -20,9 +20,14 @@ import { isDefined } from '@core/utils'; import { MobileActionResult, WidgetMobileActionResult, WidgetMobileActionType } from '@shared/models/widget.models'; import { from, of } from 'rxjs'; import { Observable } from 'rxjs/internal/Observable'; -import { catchError } from 'rxjs/operators'; +import { catchError, tap } from 'rxjs/operators'; +import { OpenDashboardMessage, ReloadUserMessage, WindowMessage } from '@shared/models/window-message.model'; +import { Params, Router } from '@angular/router'; +import { AuthService } from '@core/auth/auth.service'; const dashboardStateNameHandler = 'tbMobileDashboardStateNameHandler'; +const dashboardLoadedHandler = 'tbMobileDashboardLoadedHandler'; +const navigationHandler = 'tbMobileNavigationHandler'; const mobileHandler = 'tbMobileHandler'; // @dynamic @@ -34,10 +39,20 @@ export class MobileService { private readonly mobileApp; private readonly mobileChannel; - constructor(@Inject(WINDOW) private window: Window) { + private readonly onWindowMessageListener = this.onWindowMessage.bind(this); + + private reloadUserObservable: Observable; + private lastDashboardId: string; + + constructor(@Inject(WINDOW) private window: Window, + private router: Router, + private authService: AuthService) { const w = (this.window as any); this.mobileChannel = w.flutter_inappwebview; this.mobileApp = isDefined(this.mobileChannel); + if (this.mobileApp) { + window.addEventListener('message', this.onWindowMessageListener); + } } public isMobileApp(): boolean { @@ -50,6 +65,12 @@ export class MobileService { } } + public onDashboardLoaded() { + if (this.mobileApp) { + this.mobileChannel.callHandler(dashboardLoadedHandler); + } + } + public handleWidgetMobileAction(type: WidgetMobileActionType, ...args: any[]): Observable> { if (this.mobileApp) { @@ -67,4 +88,82 @@ export class MobileService { } } + public handleMobileNavigation(path?: string, params?: Params) { + if (this.mobileApp) { + this.mobileChannel.callHandler(navigationHandler, path, params); + } + } + + private onWindowMessage(event: MessageEvent) { + if (event.data) { + let message: WindowMessage; + try { + message = JSON.parse(event.data); + } catch (e) {} + if (message && message.type) { + switch (message.type) { + case 'openDashboardMessage': + const openDashboardMessage: OpenDashboardMessage = message.data; + this.openDashboard(openDashboardMessage); + break; + case 'reloadUserMessage': + const reloadUserMessage: ReloadUserMessage = message.data; + this.reloadUser(reloadUserMessage); + break; + } + } + } + } + + private openDashboard(openDashboardMessage: OpenDashboardMessage) { + if (openDashboardMessage && openDashboardMessage.dashboardId) { + if (this.reloadUserObservable) { + this.reloadUserObservable.subscribe( + (authenticated) => { + if (authenticated) { + this.doDashboardNavigation(openDashboardMessage); + } + } + ); + } else { + this.doDashboardNavigation(openDashboardMessage); + } + } + } + + private doDashboardNavigation(openDashboardMessage: OpenDashboardMessage) { + let url = `/dashboard/${openDashboardMessage.dashboardId}`; + const params = []; + if (openDashboardMessage.state) { + params.push(`state=${openDashboardMessage.state}`); + } + if (openDashboardMessage.embedded) { + params.push(`embedded=true`); + } + if (openDashboardMessage.hideToolbar) { + params.push(`hideToolbar=true`); + } + if (this.lastDashboardId === openDashboardMessage.dashboardId) { + params.push(`reload=${new Date().getTime()}`); + } + if (params.length) { + url += `?${params.join('&')}`; + } + this.lastDashboardId = openDashboardMessage.dashboardId; + this.router.navigateByUrl(url, {replaceUrl: true}); + } + + private reloadUser(reloadUserMessage: ReloadUserMessage) { + if (reloadUserMessage && reloadUserMessage.accessToken && reloadUserMessage.refreshToken) { + this.reloadUserObservable = this.authService.setUserFromJwtToken(reloadUserMessage.accessToken, + reloadUserMessage.refreshToken, true).pipe( + tap( + () => { + this.reloadUserObservable = null; + } + ) + ); + } + } + } diff --git a/ui-ngx/src/app/core/utils.ts b/ui-ngx/src/app/core/utils.ts index 6291438a58..a287802ec2 100644 --- a/ui-ngx/src/app/core/utils.ts +++ b/ui-ngx/src/app/core/utils.ts @@ -441,3 +441,7 @@ export function generateSecret(length?: number): string { export function validateEntityId(entityId: EntityId | null): boolean { return isDefinedAndNotNull(entityId?.id) && entityId.id !== NULL_UUID && isDefinedAndNotNull(entityId?.entityType); } + +export function isMobileApp(): boolean { + return isDefined((window as any).flutter_inappwebview); +} diff --git a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html index 5c54fe23e0..8ff6944856 100644 --- a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html +++ b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html @@ -87,7 +87,7 @@ (click)="updateDashboardImage($event)"> wallpaper -