diff --git a/application/src/main/java/org/thingsboard/server/controller/AssetController.java b/application/src/main/java/org/thingsboard/server/controller/AssetController.java index 4854c96941..d63674b7ee 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AssetController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AssetController.java @@ -16,9 +16,12 @@ package org.thingsboard.server.controller; import com.google.common.util.concurrent.ListenableFuture; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; import org.springframework.http.HttpStatus; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; @@ -34,7 +37,6 @@ import org.thingsboard.server.common.data.asset.AssetInfo; import org.thingsboard.server.common.data.asset.AssetSearchQuery; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; -import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.exception.ThingsboardException; @@ -48,6 +50,9 @@ import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.asset.AssetBulkImportService; +import org.thingsboard.server.service.importing.BulkImportRequest; +import org.thingsboard.server.service.importing.BulkImportResult; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; @@ -63,7 +68,10 @@ import static org.thingsboard.server.controller.EdgeController.EDGE_ID; @RestController @TbCoreComponent @RequestMapping("/api") +@RequiredArgsConstructor +@Slf4j public class AssetController extends BaseController { + private final AssetBulkImportService assetBulkImportService; public static final String ASSET_ID = "assetId"; @@ -108,13 +116,7 @@ public class AssetController extends BaseController { Asset savedAsset = checkNotNull(assetService.saveAsset(asset)); - logEntityAction(savedAsset.getId(), savedAsset, - savedAsset.getCustomerId(), - asset.getId() == null ? ActionType.ADDED : ActionType.UPDATED, null); - - if (asset.getId() != null) { - sendEntityNotificationMsg(savedAsset.getTenantId(), savedAsset.getId(), EdgeEventActionType.UPDATED); - } + onAssetCreatedOrUpdated(savedAsset, asset.getId() != null); return savedAsset; } catch (Exception e) { @@ -124,6 +126,20 @@ public class AssetController extends BaseController { } } + private void onAssetCreatedOrUpdated(Asset asset, boolean updated) { + try { + logEntityAction(asset.getId(), asset, + asset.getCustomerId(), + updated ? ActionType.UPDATED : ActionType.ADDED, null); + } catch (ThingsboardException e) { + log.error("Failed to log entity action", e); + } + + if (updated) { + sendEntityNotificationMsg(asset.getTenantId(), asset.getId(), EdgeEventActionType.UPDATED); + } + } + @PreAuthorize("hasAuthority('TENANT_ADMIN')") @RequestMapping(value = "/asset/{assetId}", method = RequestMethod.DELETE) @ResponseStatus(value = HttpStatus.OK) @@ -258,7 +274,7 @@ public class AssetController extends BaseController { try { TenantId tenantId = getCurrentUser().getTenantId(); PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); - if (type != null && type.trim().length()>0) { + if (type != null && type.trim().length() > 0) { return checkNotNull(assetService.findAssetsByTenantIdAndType(tenantId, type, pageLink)); } else { return checkNotNull(assetService.findAssetsByTenantId(tenantId, pageLink)); @@ -321,7 +337,7 @@ public class AssetController extends BaseController { CustomerId customerId = new CustomerId(toUUID(strCustomerId)); checkCustomerId(customerId, Operation.READ); PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); - if (type != null && type.trim().length()>0) { + if (type != null && type.trim().length() > 0) { return checkNotNull(assetService.findAssetsByTenantIdAndCustomerIdAndType(tenantId, customerId, type, pageLink)); } else { return checkNotNull(assetService.findAssetsByTenantIdAndCustomerId(tenantId, customerId, pageLink)); @@ -426,7 +442,7 @@ public class AssetController extends BaseController { @RequestMapping(value = "/edge/{edgeId}/asset/{assetId}", method = RequestMethod.POST) @ResponseBody public Asset assignAssetToEdge(@PathVariable(EDGE_ID) String strEdgeId, - @PathVariable(ASSET_ID) String strAssetId) throws ThingsboardException { + @PathVariable(ASSET_ID) String strAssetId) throws ThingsboardException { checkParameter(EDGE_ID, strEdgeId); checkParameter(ASSET_ID, strAssetId); try { @@ -444,7 +460,7 @@ public class AssetController extends BaseController { sendEntityAssignToEdgeNotificationMsg(getTenantId(), edgeId, savedAsset.getId(), EdgeEventActionType.ASSIGNED_TO_EDGE); - return savedAsset; + return savedAsset; } catch (Exception e) { logEntityAction(emptyId(EntityType.ASSET), null, @@ -530,4 +546,13 @@ public class AssetController extends BaseController { throw handleException(e); } } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @PostMapping("/asset/bulk_import") + public BulkImportResult processAssetsBulkImport(@RequestBody BulkImportRequest request) throws Exception { + return assetBulkImportService.processBulkImport(request, getCurrentUser(), importedAssetInfo -> { + onAssetCreatedOrUpdated(importedAssetInfo.getEntity(), importedAssetInfo.isUpdated()); + }); + } + } 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 c37587f546..f320b8eafa 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -121,7 +121,7 @@ import org.thingsboard.server.exception.ThingsboardErrorResponseHandler; 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.action.EntityActionService; import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.edge.EdgeLicenseService; import org.thingsboard.server.service.edge.EdgeNotificationService; @@ -279,7 +279,7 @@ public abstract class BaseController { protected EdgeLicenseService edgeLicenseService; @Autowired - protected RuleEngineEntityActionService ruleEngineEntityActionService; + protected EntityActionService entityActionService; @Value("${server.log_controller_error_stack_trace}") @Getter @@ -819,13 +819,7 @@ public abstract class BaseController { protected void logEntityAction(User user, I entityId, E entity, CustomerId customerId, ActionType actionType, Exception e, Object... additionalInfo) throws ThingsboardException { - if (customerId == null || customerId.isNullUid()) { - customerId = user.getCustomerId(); - } - if (e == null) { - ruleEngineEntityActionService.pushEntityActionToRuleEngine(entityId, entity, user.getTenantId(), customerId, actionType, user, additionalInfo); - } - auditLogService.logEntityAction(user.getTenantId(), customerId, user.getId(), user.getName(), entityId, entity, actionType, e, additionalInfo); + entityActionService.logEntityAction(user, entityId, entity, customerId, actionType, e, additionalInfo); } diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java index 05385ee948..8cfa5ee2cf 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java @@ -19,10 +19,13 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; @@ -56,7 +59,6 @@ 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.TimePageLink; -import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; @@ -67,6 +69,9 @@ import org.thingsboard.server.dao.device.claim.ReclaimResult; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.device.DeviceBulkImportService; +import org.thingsboard.server.service.importing.BulkImportRequest; +import org.thingsboard.server.service.importing.BulkImportResult; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; @@ -83,7 +88,10 @@ import static org.thingsboard.server.controller.EdgeController.EDGE_ID; @RestController @TbCoreComponent @RequestMapping("/api") +@RequiredArgsConstructor +@Slf4j public class DeviceController extends BaseController { + private final DeviceBulkImportService deviceBulkImportService; private static final String DEVICE_ID = "deviceId"; private static final String DEVICE_NAME = "deviceName"; @@ -133,11 +141,7 @@ public class DeviceController extends BaseController { Device savedDevice = checkNotNull(deviceService.saveDeviceWithAccessToken(device, accessToken)); - tbClusterService.onDeviceUpdated(savedDevice, oldDevice); - - logEntityAction(savedDevice.getId(), savedDevice, - savedDevice.getCustomerId(), - created ? ActionType.ADDED : ActionType.UPDATED, null); + onDeviceCreatedOrUpdated(savedDevice, oldDevice, !created); return savedDevice; } catch (Exception e) { @@ -148,6 +152,18 @@ public class DeviceController extends BaseController { } + private void onDeviceCreatedOrUpdated(Device savedDevice, Device oldDevice, boolean updated) { + tbClusterService.onDeviceUpdated(savedDevice, oldDevice); + + try { + logEntityAction(savedDevice.getId(), savedDevice, + savedDevice.getCustomerId(), + updated ? ActionType.UPDATED : ActionType.ADDED, null); + } catch (ThingsboardException e) { + log.error("Failed to log entity action", e); + } + } + @PreAuthorize("hasAuthority('TENANT_ADMIN')") @RequestMapping(value = "/device/{deviceId}", method = RequestMethod.DELETE) @ResponseStatus(value = HttpStatus.OK) @@ -776,4 +792,13 @@ public class DeviceController extends BaseController { throw handleException(e); } } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @PostMapping("/device/bulk_import") + public BulkImportResult processDevicesBulkImport(@RequestBody BulkImportRequest request) throws Exception { + return deviceBulkImportService.processBulkImport(request, getCurrentUser(), importedDeviceInfo -> { + onDeviceCreatedOrUpdated(importedDeviceInfo.getEntity(), importedDeviceInfo.getOldEntity(), importedDeviceInfo.isUpdated()); + }); + } + } diff --git a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java index a5855a22be..21c39cdd58 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java @@ -17,11 +17,13 @@ package org.thingsboard.server.controller; import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.ListenableFuture; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; @@ -52,10 +54,14 @@ import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.edge.EdgeBulkImportService; +import org.thingsboard.server.service.importing.BulkImportRequest; +import org.thingsboard.server.service.importing.BulkImportResult; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; +import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.stream.Collectors; @@ -64,7 +70,9 @@ import java.util.stream.Collectors; @TbCoreComponent @Slf4j @RequestMapping("/api") +@RequiredArgsConstructor public class EdgeController extends BaseController { + private final EdgeBulkImportService edgeBulkImportService; public static final String EDGE_ID = "edgeId"; @@ -132,17 +140,8 @@ public class EdgeController extends BaseController { edge.getId(), edge); Edge savedEdge = checkNotNull(edgeService.saveEdge(edge, true)); + onEdgeCreatedOrUpdated(tenantId, savedEdge, edgeTemplateRootRuleChain, !created); - if (created) { - ruleChainService.assignRuleChainToEdge(tenantId, edgeTemplateRootRuleChain.getId(), savedEdge.getId()); - edgeNotificationService.setEdgeRootRuleChain(tenantId, savedEdge, edgeTemplateRootRuleChain.getId()); - edgeService.assignDefaultRuleChainsToEdge(tenantId, savedEdge.getId()); - } - - tbClusterService.broadcastEntityStateChangeEvent(savedEdge.getTenantId(), savedEdge.getId(), - created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - - logEntityAction(savedEdge.getId(), savedEdge, null, created ? ActionType.ADDED : ActionType.UPDATED, null); return savedEdge; } catch (Exception e) { logEntityAction(emptyId(EntityType.EDGE), edge, @@ -151,6 +150,19 @@ public class EdgeController extends BaseController { } } + private void onEdgeCreatedOrUpdated(TenantId tenantId, Edge edge, RuleChain edgeTemplateRootRuleChain, boolean updated) throws IOException, ThingsboardException { + if (!updated) { + ruleChainService.assignRuleChainToEdge(tenantId, edgeTemplateRootRuleChain.getId(), edge.getId()); + edgeNotificationService.setEdgeRootRuleChain(tenantId, edge, edgeTemplateRootRuleChain.getId()); + edgeService.assignDefaultRuleChainsToEdge(tenantId, edge.getId()); + } + + tbClusterService.broadcastEntityStateChangeEvent(edge.getTenantId(), edge.getId(), + updated ? ComponentLifecycleEvent.UPDATED : ComponentLifecycleEvent.CREATED); + + logEntityAction(edge.getId(), edge, null, updated ? ActionType.UPDATED : ActionType.ADDED, null); + } + @PreAuthorize("hasAuthority('TENANT_ADMIN')") @RequestMapping(value = "/edge/{edgeId}", method = RequestMethod.DELETE) @ResponseStatus(value = HttpStatus.OK) @@ -563,6 +575,24 @@ public class EdgeController extends BaseController { } } + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @PostMapping("/edge/bulk_import") + public BulkImportResult processEdgeBulkImport(@RequestBody BulkImportRequest request) throws Exception { + SecurityUser user = getCurrentUser(); + RuleChain edgeTemplateRootRuleChain = ruleChainService.getEdgeTemplateRootRuleChain(user.getTenantId()); + if (edgeTemplateRootRuleChain == null) { + throw new DataValidationException("Root edge rule chain is not available!"); + } + + return edgeBulkImportService.processBulkImport(request, user, importedAssetInfo -> { + try { + onEdgeCreatedOrUpdated(user.getTenantId(), importedAssetInfo.getEntity(), edgeTemplateRootRuleChain, importedAssetInfo.isUpdated()); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + } + private void cleanUpLicenseKey(Edge edge) { edge.setEdgeLicenseKey(null); } diff --git a/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java b/application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java similarity index 93% rename from application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java rename to application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java index ef7d077edb..27aef2d42e 100644 --- a/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java +++ b/application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java @@ -28,6 +28,7 @@ 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.exception.ThingsboardException; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -38,6 +39,7 @@ 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.dao.audit.AuditLogService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.cluster.TbClusterService; @@ -49,8 +51,9 @@ import java.util.stream.Collectors; @Service @RequiredArgsConstructor @Slf4j -public class RuleEngineEntityActionService { +public class EntityActionService { private final TbClusterService tbClusterService; + private final AuditLogService auditLogService; private static final ObjectMapper json = new ObjectMapper(); @@ -209,6 +212,17 @@ public class RuleEngineEntityActionService { } } + public void logEntityAction(User user, I entityId, E entity, CustomerId customerId, + ActionType actionType, Exception e, Object... additionalInfo) { + if (customerId == null || customerId.isNullUid()) { + customerId = user.getCustomerId(); + } + if (e == null) { + pushEntityActionToRuleEngine(entityId, entity, user.getTenantId(), customerId, actionType, user, additionalInfo); + } + auditLogService.logEntityAction(user.getTenantId(), customerId, user.getId(), user.getName(), entityId, entity, actionType, e, additionalInfo); + } + private T extractParameter(Class clazz, int index, Object... additionalInfo) { T result = null; diff --git a/application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java new file mode 100644 index 0000000000..2f94ac3b0d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java @@ -0,0 +1,94 @@ +/** + * 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.asset; + +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.fasterxml.jackson.databind.node.TextNode; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.action.EntityActionService; +import org.thingsboard.server.service.importing.AbstractBulkImportService; +import org.thingsboard.server.service.importing.BulkImportColumnType; +import org.thingsboard.server.service.importing.BulkImportRequest; +import org.thingsboard.server.service.importing.ImportedEntityInfo; +import org.thingsboard.server.service.security.AccessValidator; +import org.thingsboard.server.service.security.model.SecurityUser; +import org.thingsboard.server.service.security.permission.AccessControlService; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; + +import java.util.Map; +import java.util.Optional; + +@Service +@TbCoreComponent +public class AssetBulkImportService extends AbstractBulkImportService { + private final AssetService assetService; + + public AssetBulkImportService(TelemetrySubscriptionService tsSubscriptionService, TbTenantProfileCache tenantProfileCache, + AccessControlService accessControlService, AccessValidator accessValidator, + EntityActionService entityActionService, TbClusterService clusterService, AssetService assetService) { + super(tsSubscriptionService, tenantProfileCache, accessControlService, accessValidator, entityActionService, clusterService); + this.assetService = assetService; + } + + @Override + protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user) { + ImportedEntityInfo importedEntityInfo = new ImportedEntityInfo<>(); + + Asset asset = new Asset(); + asset.setTenantId(user.getTenantId()); + setAssetFields(asset, fields); + + Asset existingAsset = assetService.findAssetByTenantIdAndName(user.getTenantId(), asset.getName()); + if (existingAsset != null && importRequest.getMapping().getUpdate()) { + importedEntityInfo.setOldEntity(new Asset(existingAsset)); + importedEntityInfo.setUpdated(true); + existingAsset.update(asset); + asset = existingAsset; + } + asset = assetService.saveAsset(asset); + + importedEntityInfo.setEntity(asset); + return importedEntityInfo; + } + + private void setAssetFields(Asset asset, Map fields) { + ObjectNode additionalInfo = (ObjectNode) Optional.ofNullable(asset.getAdditionalInfo()).orElseGet(JacksonUtil::newObjectNode); + fields.forEach((columnType, value) -> { + switch (columnType) { + case NAME: + asset.setName(value); + break; + case TYPE: + asset.setType(value); + break; + case LABEL: + asset.setLabel(value); + break; + case DESCRIPTION: + additionalInfo.set("description", new TextNode(value)); + break; + } + }); + asset.setAdditionalInfo(additionalInfo); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java new file mode 100644 index 0000000000..82d4dc5571 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java @@ -0,0 +1,259 @@ +/** + * 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.device; + +import com.fasterxml.jackson.databind.node.BooleanNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.fasterxml.jackson.databind.node.TextNode; +import lombok.SneakyThrows; +import org.apache.commons.collections.CollectionUtils; +import org.apache.commons.lang3.RandomStringUtils; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.DeviceProfileProvisionType; +import org.thingsboard.server.common.data.DeviceProfileType; +import org.thingsboard.server.common.data.DeviceTransportType; +import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MClientCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode; +import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileConfiguration; +import org.thingsboard.server.common.data.device.profile.DeviceProfileData; +import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; +import org.thingsboard.server.common.data.device.profile.DisabledDeviceProfileProvisionConfiguration; +import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.data.security.DeviceCredentialsType; +import org.thingsboard.server.dao.device.DeviceCredentialsService; +import org.thingsboard.server.dao.device.DeviceProfileService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.action.EntityActionService; +import org.thingsboard.server.service.importing.AbstractBulkImportService; +import org.thingsboard.server.service.importing.BulkImportColumnType; +import org.thingsboard.server.service.importing.BulkImportRequest; +import org.thingsboard.server.service.importing.ImportedEntityInfo; +import org.thingsboard.server.service.security.AccessValidator; +import org.thingsboard.server.service.security.model.SecurityUser; +import org.thingsboard.server.service.security.permission.AccessControlService; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; + +import java.util.Collection; +import java.util.EnumSet; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +@Service +@TbCoreComponent +public class DeviceBulkImportService extends AbstractBulkImportService { + protected final DeviceService deviceService; + protected final DeviceCredentialsService deviceCredentialsService; + protected final DeviceProfileService deviceProfileService; + + public DeviceBulkImportService(TelemetrySubscriptionService tsSubscriptionService, TbTenantProfileCache tenantProfileCache, + AccessControlService accessControlService, AccessValidator accessValidator, + EntityActionService entityActionService, TbClusterService clusterService, + DeviceService deviceService, DeviceCredentialsService deviceCredentialsService, + DeviceProfileService deviceProfileService) { + super(tsSubscriptionService, tenantProfileCache, accessControlService, accessValidator, entityActionService, clusterService); + this.deviceService = deviceService; + this.deviceCredentialsService = deviceCredentialsService; + this.deviceProfileService = deviceProfileService; + } + + @Override + protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user) { + ImportedEntityInfo importedEntityInfo = new ImportedEntityInfo<>(); + + Device device = new Device(); + device.setTenantId(user.getTenantId()); + setDeviceFields(device, fields); + + Device existingDevice = deviceService.findDeviceByTenantIdAndName(user.getTenantId(), device.getName()); + if (existingDevice != null && importRequest.getMapping().getUpdate()) { + importedEntityInfo.setOldEntity(new Device(existingDevice)); + importedEntityInfo.setUpdated(true); + existingDevice.updateDevice(device); + device = existingDevice; + } + + DeviceCredentials deviceCredentials; + try { + deviceCredentials = createDeviceCredentials(fields); + deviceCredentialsService.formatCredentials(deviceCredentials); + } catch (Exception e) { + throw new DeviceCredentialsValidationException("Invalid device credentials: " + e.getMessage()); + } + + if (deviceCredentials.getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) { + setUpLwM2mDeviceProfile(user.getTenantId(), device); + } + + device = deviceService.saveDeviceWithCredentials(device, deviceCredentials); + + importedEntityInfo.setEntity(device); + return importedEntityInfo; + } + + private void setDeviceFields(Device device, Map fields) { + ObjectNode additionalInfo = (ObjectNode) Optional.ofNullable(device.getAdditionalInfo()).orElseGet(JacksonUtil::newObjectNode); + fields.forEach((columnType, value) -> { + switch (columnType) { + case NAME: + device.setName(value); + break; + case TYPE: + device.setType(value); + break; + case LABEL: + device.setLabel(value); + break; + case DESCRIPTION: + additionalInfo.set("description", new TextNode(value)); + break; + case IS_GATEWAY: + additionalInfo.set("gateway", BooleanNode.valueOf(Boolean.parseBoolean(value))); + break; + } + device.setAdditionalInfo(additionalInfo); + }); + } + + @SneakyThrows + private DeviceCredentials createDeviceCredentials(Map fields) { + DeviceCredentials credentials = new DeviceCredentials(); + if (fields.containsKey(BulkImportColumnType.LWM2M_CLIENT_ENDPOINT)) { + credentials.setCredentialsType(DeviceCredentialsType.LWM2M_CREDENTIALS); + setUpLwm2mCredentials(fields, credentials); + } else if (fields.containsKey(BulkImportColumnType.X509)) { + credentials.setCredentialsType(DeviceCredentialsType.X509_CERTIFICATE); + setUpX509CertificateCredentials(fields, credentials); + } else if (CollectionUtils.containsAny(fields.keySet(), EnumSet.of(BulkImportColumnType.MQTT_CLIENT_ID, BulkImportColumnType.MQTT_USER_NAME, BulkImportColumnType.MQTT_PASSWORD))) { + credentials.setCredentialsType(DeviceCredentialsType.MQTT_BASIC); + setUpBasicMqttCredentials(fields, credentials); + } else { + credentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); + setUpAccessTokenCredentials(fields, credentials); + } + return credentials; + } + + private void setUpAccessTokenCredentials(Map fields, DeviceCredentials credentials) { + credentials.setCredentialsId(Optional.ofNullable(fields.get(BulkImportColumnType.ACCESS_TOKEN)) + .orElseGet(() -> RandomStringUtils.randomAlphanumeric(20))); + } + + private void setUpBasicMqttCredentials(Map fields, DeviceCredentials credentials) { + BasicMqttCredentials basicMqttCredentials = new BasicMqttCredentials(); + basicMqttCredentials.setClientId(fields.get(BulkImportColumnType.MQTT_CLIENT_ID)); + basicMqttCredentials.setUserName(fields.get(BulkImportColumnType.MQTT_USER_NAME)); + basicMqttCredentials.setPassword(fields.get(BulkImportColumnType.MQTT_PASSWORD)); + credentials.setCredentialsValue(JacksonUtil.toString(basicMqttCredentials)); + } + + private void setUpX509CertificateCredentials(Map fields, DeviceCredentials credentials) { + credentials.setCredentialsValue(fields.get(BulkImportColumnType.X509)); + } + + private void setUpLwm2mCredentials(Map fields, DeviceCredentials credentials) throws com.fasterxml.jackson.core.JsonProcessingException { + ObjectNode lwm2mCredentials = JacksonUtil.newObjectNode(); + + Set.of(BulkImportColumnType.LWM2M_CLIENT_SECURITY_CONFIG_MODE, BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_SECURITY_MODE, + BulkImportColumnType.LWM2M_SERVER_SECURITY_MODE).stream() + .map(fields::get) + .filter(Objects::nonNull) + .forEach(securityMode -> { + try { + LwM2MSecurityMode.valueOf(securityMode.toUpperCase()); + } catch (IllegalArgumentException e) { + throw new DeviceCredentialsValidationException("Unknown LwM2M security mode: " + securityMode + ", (the mode should be: NO_SEC, PSK, RPK, X509)!"); + } + }); + + ObjectNode client = JacksonUtil.newObjectNode(); + setValues(client, fields, Set.of(BulkImportColumnType.LWM2M_CLIENT_SECURITY_CONFIG_MODE, + BulkImportColumnType.LWM2M_CLIENT_ENDPOINT, BulkImportColumnType.LWM2M_CLIENT_IDENTITY, + BulkImportColumnType.LWM2M_CLIENT_KEY, BulkImportColumnType.LWM2M_CLIENT_CERT)); + LwM2MClientCredentials lwM2MClientCredentials = JacksonUtil.treeToValue(client, LwM2MClientCredentials.class); + // so that only fields needed for specific type of lwM2MClientCredentials were saved in json + lwm2mCredentials.set("client", JacksonUtil.valueToTree(lwM2MClientCredentials)); + + ObjectNode bootstrapServer = JacksonUtil.newObjectNode(); + setValues(bootstrapServer, fields, Set.of(BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_SECURITY_MODE, + BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_PUBLIC_KEY_OR_ID, BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_SECRET_KEY)); + + ObjectNode lwm2mServer = JacksonUtil.newObjectNode(); + setValues(lwm2mServer, fields, Set.of(BulkImportColumnType.LWM2M_SERVER_SECURITY_MODE, + BulkImportColumnType.LWM2M_SERVER_CLIENT_PUBLIC_KEY_OR_ID, BulkImportColumnType.LWM2M_SERVER_CLIENT_SECRET_KEY)); + + ObjectNode bootstrap = JacksonUtil.newObjectNode(); + bootstrap.set("bootstrapServer", bootstrapServer); + bootstrap.set("lwm2mServer", lwm2mServer); + lwm2mCredentials.set("bootstrap", bootstrap); + + credentials.setCredentialsValue(lwm2mCredentials.toString()); + } + + private void setUpLwM2mDeviceProfile(TenantId tenantId, Device device) { + DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileByName(tenantId, device.getType()); + if (deviceProfile != null) { + if (deviceProfile.getTransportType() != DeviceTransportType.LWM2M) { + deviceProfile.setTransportType(DeviceTransportType.LWM2M); + deviceProfile.getProfileData().setTransportConfiguration(new Lwm2mDeviceProfileTransportConfiguration()); + deviceProfile = deviceProfileService.saveDeviceProfile(deviceProfile); + device.setDeviceProfileId(deviceProfile.getId()); + } + } else { + deviceProfile = new DeviceProfile(); + deviceProfile.setTenantId(tenantId); + deviceProfile.setType(DeviceProfileType.DEFAULT); + deviceProfile.setName(device.getType()); + deviceProfile.setTransportType(DeviceTransportType.LWM2M); + deviceProfile.setProvisionType(DeviceProfileProvisionType.DISABLED); + + DeviceProfileData deviceProfileData = new DeviceProfileData(); + DefaultDeviceProfileConfiguration configuration = new DefaultDeviceProfileConfiguration(); + DeviceProfileTransportConfiguration transportConfiguration = new Lwm2mDeviceProfileTransportConfiguration(); + DisabledDeviceProfileProvisionConfiguration provisionConfiguration = new DisabledDeviceProfileProvisionConfiguration(null); + + deviceProfileData.setConfiguration(configuration); + deviceProfileData.setTransportConfiguration(transportConfiguration); + deviceProfileData.setProvisionConfiguration(provisionConfiguration); + deviceProfile.setProfileData(deviceProfileData); + + deviceProfile = deviceProfileService.saveDeviceProfile(deviceProfile); + device.setDeviceProfileId(deviceProfile.getId()); + } + } + + private void setValues(ObjectNode objectNode, Map data, Collection columns) { + for (BulkImportColumnType column : columns) { + String value = StringUtils.defaultString(data.get(column), column.getDefaultValue()); + if (value != null && column.getKey() != null) { + objectNode.set(column.getKey(), new TextNode(value)); + } + } + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java new file mode 100644 index 0000000000..ec6a2a4e55 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java @@ -0,0 +1,106 @@ +/** + * 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.edge; + +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.fasterxml.jackson.databind.node.TextNode; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.dao.edge.EdgeService; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.action.EntityActionService; +import org.thingsboard.server.service.importing.AbstractBulkImportService; +import org.thingsboard.server.service.importing.BulkImportColumnType; +import org.thingsboard.server.service.importing.BulkImportRequest; +import org.thingsboard.server.service.importing.ImportedEntityInfo; +import org.thingsboard.server.service.security.AccessValidator; +import org.thingsboard.server.service.security.model.SecurityUser; +import org.thingsboard.server.service.security.permission.AccessControlService; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; + +import java.util.Map; +import java.util.Optional; + +@Service +@TbCoreComponent +public class EdgeBulkImportService extends AbstractBulkImportService { + private final EdgeService edgeService; + + public EdgeBulkImportService(TelemetrySubscriptionService tsSubscriptionService, TbTenantProfileCache tenantProfileCache, + AccessControlService accessControlService, AccessValidator accessValidator, + EntityActionService entityActionService, TbClusterService clusterService, EdgeService edgeService) { + super(tsSubscriptionService, tenantProfileCache, accessControlService, accessValidator, entityActionService, clusterService); + this.edgeService = edgeService; + } + + @Override + protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user) { + ImportedEntityInfo importedEntityInfo = new ImportedEntityInfo<>(); + + Edge edge = new Edge(); + edge.setTenantId(user.getTenantId()); + setEdgeFields(edge, fields); + + Edge existingEdge = edgeService.findEdgeByTenantIdAndName(user.getTenantId(), edge.getName()); + if (existingEdge != null && importRequest.getMapping().getUpdate()) { + importedEntityInfo.setOldEntity(new Edge(existingEdge)); + importedEntityInfo.setUpdated(true); + existingEdge.update(edge); + edge = existingEdge; + } + edge = edgeService.saveEdge(edge, true); + + importedEntityInfo.setEntity(edge); + return importedEntityInfo; + } + + private void setEdgeFields(Edge edge, Map fields) { + ObjectNode additionalInfo = (ObjectNode) Optional.ofNullable(edge.getAdditionalInfo()).orElseGet(JacksonUtil::newObjectNode); + fields.forEach((columnType, value) -> { + switch (columnType) { + case NAME: + edge.setName(value); + break; + case TYPE: + edge.setType(value); + break; + case LABEL: + edge.setLabel(value); + break; + case DESCRIPTION: + additionalInfo.set("description", new TextNode(value)); + break; + case EDGE_LICENSE_KEY: + edge.setEdgeLicenseKey(value); + break; + case CLOUD_ENDPOINT: + edge.setCloudEndpoint(value); + break; + case ROUTING_KEY: + edge.setRoutingKey(value); + break; + case SECRET: + edge.setSecret(value); + break; + } + }); + edge.setAdditionalInfo(additionalInfo); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java new file mode 100644 index 0000000000..b1d0b30c9d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java @@ -0,0 +1,245 @@ +/** + * 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.importing; + +import com.google.common.util.concurrent.FutureCallback; +import com.google.gson.JsonObject; +import com.google.gson.JsonPrimitive; +import lombok.Data; +import lombok.RequiredArgsConstructor; +import lombok.SneakyThrows; +import org.apache.commons.lang3.StringUtils; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.BaseData; +import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.audit.ActionType; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.UUIDBased; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.DataType; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.transport.adaptor.JsonConverter; +import org.thingsboard.server.controller.BaseController; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import org.thingsboard.server.service.action.EntityActionService; +import org.thingsboard.server.service.importing.BulkImportRequest.ColumnMapping; +import org.thingsboard.server.service.security.AccessValidator; +import org.thingsboard.server.service.security.model.SecurityUser; +import org.thingsboard.server.service.security.permission.AccessControlService; +import org.thingsboard.server.service.security.permission.Operation; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; +import org.thingsboard.server.utils.CsvUtils; +import org.thingsboard.server.utils.TypeCastUtil; + +import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Consumer; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +@RequiredArgsConstructor +public abstract class AbstractBulkImportService> { + protected final TelemetrySubscriptionService tsSubscriptionService; + protected final TbTenantProfileCache tenantProfileCache; + protected final AccessControlService accessControlService; + protected final AccessValidator accessValidator; + protected final EntityActionService entityActionService; + protected final TbClusterService clusterService; + + public final BulkImportResult processBulkImport(BulkImportRequest request, SecurityUser user, Consumer> onEntityImported) throws Exception { + BulkImportResult result = new BulkImportResult<>(); + + AtomicInteger i = new AtomicInteger(0); + if (request.getMapping().getHeader()) { + i.incrementAndGet(); + } + + parseData(request).forEach(entityData -> { + i.incrementAndGet(); + try { + ImportedEntityInfo importedEntityInfo = saveEntity(request, entityData.getFields(), user); + onEntityImported.accept(importedEntityInfo); + + E entity = importedEntityInfo.getEntity(); + + saveKvs(user, entity, entityData.getKvs()); + + if (importedEntityInfo.getRelatedError() != null) { + throw new RuntimeException(importedEntityInfo.getRelatedError()); + } + + if (importedEntityInfo.isUpdated()) { + result.setUpdated(result.getUpdated() + 1); + } else { + result.setCreated(result.getCreated() + 1); + } + } catch (Exception e) { + result.setErrors(result.getErrors() + 1); + result.getErrorsList().add(String.format("Line %d: %s", i.get(), e.getMessage())); + } + }); + + return result; + } + + protected abstract ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user); + + /* + * Attributes' values are firstly added to JsonObject in order to then make some type cast, + * because we get all values as strings from CSV + * */ + private void saveKvs(SecurityUser user, E entity, Map data) { + Arrays.stream(BulkImportColumnType.values()) + .filter(BulkImportColumnType::isKv) + .map(kvType -> { + JsonObject kvs = new JsonObject(); + data.entrySet().stream() + .filter(dataEntry -> dataEntry.getKey().getType() == kvType && + StringUtils.isNotEmpty(dataEntry.getKey().getKey())) + .forEach(dataEntry -> kvs.add(dataEntry.getKey().getKey(), dataEntry.getValue().toJsonPrimitive())); + return Map.entry(kvType, kvs); + }) + .filter(kvsEntry -> kvsEntry.getValue().entrySet().size() > 0) + .forEach(kvsEntry -> { + BulkImportColumnType kvType = kvsEntry.getKey(); + if (kvType == BulkImportColumnType.SHARED_ATTRIBUTE || kvType == BulkImportColumnType.SERVER_ATTRIBUTE) { + saveAttributes(user, entity, kvsEntry, kvType); + } else { + saveTelemetry(user, entity, kvsEntry); + } + }); + } + + @SneakyThrows + private void saveTelemetry(SecurityUser user, E entity, Map.Entry kvsEntry) { + List timeseries = JsonConverter.convertToTelemetry(kvsEntry.getValue(), System.currentTimeMillis()) + .entrySet().stream() + .flatMap(entry -> entry.getValue().stream().map(kvEntry -> new BasicTsKvEntry(entry.getKey(), kvEntry))) + .collect(Collectors.toList()); + + accessValidator.validateEntityAndCallback(user, Operation.WRITE_TELEMETRY, entity.getId(), (result, tenantId, entityId) -> { + TenantProfile tenantProfile = tenantProfileCache.get(tenantId); + long tenantTtl = TimeUnit.DAYS.toSeconds(((DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration()).getDefaultStorageTtlDays()); + tsSubscriptionService.saveAndNotify(tenantId, user.getCustomerId(), entityId, timeseries, tenantTtl, new FutureCallback() { + @Override + public void onSuccess(@Nullable Void tmp) { + entityActionService.logEntityAction(user, (UUIDBased & EntityId) entityId, null, null, + ActionType.TIMESERIES_UPDATED, null, timeseries); + } + + @Override + public void onFailure(Throwable t) { + entityActionService.logEntityAction(user, (UUIDBased & EntityId) entityId, null, null, + ActionType.TIMESERIES_UPDATED, BaseController.toException(t), timeseries); + throw new RuntimeException(t); + } + }); + }); + } + + @SneakyThrows + private void saveAttributes(SecurityUser user, E entity, Map.Entry kvsEntry, BulkImportColumnType kvType) { + String scope = kvType.getKey(); + List attributes = new ArrayList<>(JsonConverter.convertToAttributes(kvsEntry.getValue())); + + accessValidator.validateEntityAndCallback(user, Operation.WRITE_ATTRIBUTES, entity.getId(), (result, tenantId, entityId) -> { + tsSubscriptionService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback<>() { + + @Override + public void onSuccess(Void unused) { + entityActionService.logEntityAction(user, (UUIDBased & EntityId) entityId, null, + null, ActionType.ATTRIBUTES_UPDATED, null, scope, attributes); + } + + @Override + public void onFailure(Throwable throwable) { + entityActionService.logEntityAction(user, (UUIDBased & EntityId) entityId, null, + null, ActionType.ATTRIBUTES_UPDATED, BaseController.toException(throwable), + scope, attributes); + throw new RuntimeException(throwable); + } + + }); + }); + } + + private List parseData(BulkImportRequest request) throws Exception { + List> records = CsvUtils.parseCsv(request.getFile(), request.getMapping().getDelimiter()); + if (request.getMapping().getHeader()) { + records.remove(0); + } + + List columnsMappings = request.getMapping().getColumns(); + return records.stream() + .map(record -> { + EntityData entityData = new EntityData(); + Stream.iterate(0, i -> i < record.size(), i -> i + 1) + .map(i -> Map.entry(columnsMappings.get(i), record.get(i))) + .filter(entry -> StringUtils.isNotEmpty(entry.getValue())) + .forEach(entry -> { + if (!entry.getKey().getType().isKv()) { + entityData.getFields().put(entry.getKey().getType(), entry.getValue()); + } else { + Map.Entry castResult = TypeCastUtil.castValue(entry.getValue()); + entityData.getKvs().put(entry.getKey(), new ParsedValue(castResult.getValue(), castResult.getKey())); + } + }); + return entityData; + }) + .collect(Collectors.toList()); + } + + @Data + protected static class EntityData { + private final Map fields = new LinkedHashMap<>(); + private final Map kvs = new LinkedHashMap<>(); + } + + @Data + protected static class ParsedValue { + private final Object value; + private final DataType dataType; + + public JsonPrimitive toJsonPrimitive() { + switch (dataType) { + case STRING: + return new JsonPrimitive((String) value); + case LONG: + return new JsonPrimitive((Long) value); + case DOUBLE: + return new JsonPrimitive((Double) value); + case BOOLEAN: + return new JsonPrimitive((Boolean) value); + default: + return null; + } + } + + public String stringValue() { + return value.toString(); + } + + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java new file mode 100644 index 0000000000..bf06753636 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java @@ -0,0 +1,77 @@ +/** + * 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.importing; + +import lombok.Getter; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode; + +@Getter +public enum BulkImportColumnType { + NAME, + TYPE, + LABEL, + SHARED_ATTRIBUTE(DataConstants.SHARED_SCOPE, true), + SERVER_ATTRIBUTE(DataConstants.SERVER_SCOPE, true), + TIMESERIES(true), + ACCESS_TOKEN, + X509, + MQTT_CLIENT_ID, + MQTT_USER_NAME, + MQTT_PASSWORD, + LWM2M_CLIENT_ENDPOINT("endpoint"), + LWM2M_CLIENT_SECURITY_CONFIG_MODE("securityConfigClientMode", LwM2MSecurityMode.NO_SEC.name()), + LWM2M_CLIENT_IDENTITY("identity"), + LWM2M_CLIENT_KEY("key"), + LWM2M_CLIENT_CERT("cert"), + LWM2M_BOOTSTRAP_SERVER_SECURITY_MODE("securityMode", LwM2MSecurityMode.NO_SEC.name()), + LWM2M_BOOTSTRAP_SERVER_PUBLIC_KEY_OR_ID("clientPublicKeyOrId"), + LWM2M_BOOTSTRAP_SERVER_SECRET_KEY("clientSecretKey"), + LWM2M_SERVER_SECURITY_MODE("securityMode", LwM2MSecurityMode.NO_SEC.name()), + LWM2M_SERVER_CLIENT_PUBLIC_KEY_OR_ID("clientPublicKeyOrId"), + LWM2M_SERVER_CLIENT_SECRET_KEY("clientSecretKey"), + IS_GATEWAY, + DESCRIPTION, + EDGE_LICENSE_KEY, + CLOUD_ENDPOINT, + ROUTING_KEY, + SECRET; + + private String key; + private String defaultValue; + private boolean isKv = false; + + BulkImportColumnType() { + } + + BulkImportColumnType(String key) { + this.key = key; + } + + BulkImportColumnType(String key, String defaultValue) { + this.key = key; + this.defaultValue = defaultValue; + } + + BulkImportColumnType(boolean isKv) { + this.isKv = isKv; + } + + BulkImportColumnType(String key, boolean isKv) { + this.key = key; + this.isKv = isKv; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/importing/BulkImportRequest.java b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportRequest.java new file mode 100644 index 0000000000..4059cc52d5 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportRequest.java @@ -0,0 +1,41 @@ +/** + * 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.importing; + +import lombok.Data; + +import java.util.List; + +@Data +public class BulkImportRequest { + private String file; + private Mapping mapping; + + @Data + public static class Mapping { + private List columns; + private Character delimiter; + private Boolean update; + private Boolean header; + } + + @Data + public static class ColumnMapping { + private BulkImportColumnType type; + private String key; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/importing/BulkImportResult.java b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportResult.java new file mode 100644 index 0000000000..d6fa6ccbf9 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportResult.java @@ -0,0 +1,30 @@ +/** + * 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.importing; + +import lombok.Data; + +import java.util.LinkedList; +import java.util.List; + +@Data +public class BulkImportResult { + private int created = 0; + private int updated = 0; + private int errors = 0; + private List errorsList = new LinkedList<>(); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/importing/ImportedEntityInfo.java b/application/src/main/java/org/thingsboard/server/service/importing/ImportedEntityInfo.java new file mode 100644 index 0000000000..958863537a --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/importing/ImportedEntityInfo.java @@ -0,0 +1,26 @@ +/** + * 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.importing; + +import lombok.Data; + +@Data +public class ImportedEntityInfo { + private E entity; + private boolean isUpdated; + private E oldEntity; + private String relatedError; +} diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java index 051a6c92b2..1dfb4a67bf 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java @@ -26,7 +26,6 @@ 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; @@ -34,17 +33,12 @@ 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 org.thingsboard.server.service.action.EntityActionService; -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 @@ -60,7 +54,7 @@ public class AlarmsCleanUpService { private final AlarmDao alarmDao; private final AlarmService alarmService; private final RelationService relationService; - private final RuleEngineEntityActionService ruleEngineEntityActionService; + private final EntityActionService entityActionService; private final PartitionService partitionService; private final TbTenantProfileCache tenantProfileCache; @@ -90,7 +84,7 @@ public class AlarmsCleanUpService { 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); + entityActionService.pushEntityActionToRuleEngine(alarm.getOriginator(), alarm, tenantId, null, ActionType.ALARM_DELETE, null); }); totalRemoved += toRemove.getTotalElements(); diff --git a/application/src/main/java/org/thingsboard/server/utils/CsvUtils.java b/application/src/main/java/org/thingsboard/server/utils/CsvUtils.java new file mode 100644 index 0000000000..d8b6dff8a7 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/utils/CsvUtils.java @@ -0,0 +1,46 @@ +/** + * 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.utils; + +import lombok.AccessLevel; +import lombok.NoArgsConstructor; +import org.apache.commons.csv.CSVFormat; +import org.apache.commons.csv.CSVRecord; +import org.apache.commons.io.input.CharSequenceReader; + +import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +@NoArgsConstructor(access = AccessLevel.PRIVATE) +public class CsvUtils { + + public static List> parseCsv(String content, Character delimiter) throws Exception { + CSVFormat csvFormat = delimiter.equals(',') ? CSVFormat.DEFAULT : CSVFormat.DEFAULT.withDelimiter(delimiter); + + List records; + try (CharSequenceReader reader = new CharSequenceReader(content)) { + records = csvFormat.parse(reader).getRecords(); + } + + return records.stream() + .map(record -> Stream.iterate(0, i -> i < record.size(), i -> i + 1) + .map(record::get) + .collect(Collectors.toList())) + .collect(Collectors.toList()); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java b/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java new file mode 100644 index 0000000000..6bc9a4d249 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java @@ -0,0 +1,55 @@ +/** + * 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.utils; + +import org.apache.commons.lang3.math.NumberUtils; +import org.thingsboard.server.common.data.kv.DataType; + +import java.math.BigDecimal; +import java.util.Map; + +public class TypeCastUtil { + + private TypeCastUtil() {} + + public static Map.Entry castValue(String value) { + if (isNumber(value)) { + String formattedValue = value.replace(',', '.'); + try { + BigDecimal bd = new BigDecimal(formattedValue); + if (bd.stripTrailingZeros().scale() > 0 || isSimpleDouble(formattedValue)) { + if (bd.scale() <= 16) { + return Map.entry(DataType.DOUBLE, bd.doubleValue()); + } + } else { + return Map.entry(DataType.LONG, bd.longValueExact()); + } + } catch (RuntimeException ignored) {} + } else if (value.equalsIgnoreCase("true") || value.equalsIgnoreCase("false")) { + return Map.entry(DataType.BOOLEAN, Boolean.parseBoolean(value)); + } + return Map.entry(DataType.STRING, value); + } + + private static boolean isNumber(String value) { + return NumberUtils.isNumber(value.replace(',', '.')); + } + + private static boolean isSimpleDouble(String valueAsString) { + return valueAsString.contains(".") && !valueAsString.contains("E") && !valueAsString.contains("e"); + } + +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsService.java index e131064953..29572bf51a 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsService.java @@ -29,5 +29,7 @@ public interface DeviceCredentialsService { DeviceCredentials createDeviceCredentials(TenantId tenantId, DeviceCredentials deviceCredentials); + void formatCredentials(DeviceCredentials deviceCredentials); + void deleteDeviceCredentials(TenantId tenantId, DeviceCredentials deviceCredentials); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/Device.java b/common/data/src/main/java/org/thingsboard/server/common/data/Device.java index 9abc619b48..c5519d329b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/Device.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/Device.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.validation.NoXss; import java.io.ByteArrayInputStream; import java.io.IOException; +import java.util.Optional; @EqualsAndHashCode(callSuper = true) @Slf4j @@ -83,6 +84,7 @@ public class Device extends SearchTextBasedWithAdditionalInfo implemen this.setDeviceData(device.getDeviceData()); this.setFirmwareId(device.getFirmwareId()); this.setSoftwareId(device.getSoftwareId()); + Optional.ofNullable(device.getAdditionalInfo()).ifPresent(this::setAdditionalInfo); return this; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/asset/Asset.java b/common/data/src/main/java/org/thingsboard/server/common/data/asset/Asset.java index f9d64cb712..594e6b635f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/asset/Asset.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/asset/Asset.java @@ -25,6 +25,8 @@ import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.validation.NoXss; +import java.util.Optional; + @EqualsAndHashCode(callSuper = true) public class Asset extends SearchTextBasedWithAdditionalInfo implements HasName, HasTenantId, HasCustomerId { @@ -56,6 +58,15 @@ public class Asset extends SearchTextBasedWithAdditionalInfo implements this.label = asset.getLabel(); } + public void update(Asset asset) { + this.tenantId = asset.getTenantId(); + this.customerId = asset.getCustomerId(); + this.name = asset.getName(); + this.type = asset.getType(); + this.label = asset.getLabel(); + Optional.ofNullable(asset.getAdditionalInfo()).ifPresent(this::setAdditionalInfo); + } + public TenantId getTenantId() { return tenantId; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/AbstractLwM2MServerCredentialsWithKeys.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/AbstractLwM2MServerCredentialsWithKeys.java new file mode 100644 index 0000000000..c45fb3891b --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/AbstractLwM2MServerCredentialsWithKeys.java @@ -0,0 +1,45 @@ +/** + * 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.device.credentials.lwm2m; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.Getter; +import lombok.Setter; +import lombok.SneakyThrows; +import org.apache.commons.codec.binary.Hex; + +@Getter +@Setter +public abstract class AbstractLwM2MServerCredentialsWithKeys implements LwM2MServerCredentials { + + private String clientPublicKeyOrId; + private String clientSecretKey; + + @JsonIgnore + public byte[] getDecodedClientPublicKeyOrId() { + return getDecoded(clientPublicKeyOrId); + } + + @JsonIgnore + public byte[] getDecodedClientSecretKey() { + return getDecoded(clientSecretKey); + } + + @SneakyThrows + private static byte[] getDecoded(String key) { + return Hex.decodeHex(key.toLowerCase().toCharArray()); + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MBootstrapCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MBootstrapCredentials.java new file mode 100644 index 0000000000..70a4846e86 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MBootstrapCredentials.java @@ -0,0 +1,26 @@ +/** + * 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.device.credentials.lwm2m; + +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +public class LwM2MBootstrapCredentials { + private LwM2MServerCredentials bootstrapServer; + private LwM2MServerCredentials lwm2mServer; +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MClientCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MClientCredentials.java index adf0c2ae62..7322c80359 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MClientCredentials.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MClientCredentials.java @@ -16,6 +16,7 @@ package org.thingsboard.server.common.data.device.credentials.lwm2m; import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; @@ -26,7 +27,9 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; @JsonSubTypes.Type(value = NoSecClientCredentials.class, name = "NO_SEC"), @JsonSubTypes.Type(value = PSKClientCredentials.class, name = "PSK"), @JsonSubTypes.Type(value = RPKClientCredentials.class, name = "RPK"), - @JsonSubTypes.Type(value = X509ClientCredentials.class, name = "X509")}) + @JsonSubTypes.Type(value = X509ClientCredentials.class, name = "X509") +}) +@JsonIgnoreProperties(ignoreUnknown = true) public interface LwM2MClientCredentials { @JsonIgnore diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MDeviceCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MDeviceCredentials.java new file mode 100644 index 0000000000..3275b4313e --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MDeviceCredentials.java @@ -0,0 +1,26 @@ +/** + * 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.device.credentials.lwm2m; + +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +public class LwM2MDeviceCredentials { + private LwM2MClientCredentials client; + private LwM2MBootstrapCredentials bootstrap; +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MServerCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MServerCredentials.java new file mode 100644 index 0000000000..82e2b6feae --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/LwM2MServerCredentials.java @@ -0,0 +1,37 @@ +/** + * 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.device.credentials.lwm2m; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; + +@JsonTypeInfo( + use = JsonTypeInfo.Id.NAME, + property = "securityMode") +@JsonSubTypes({ + @JsonSubTypes.Type(value = NoSecServerCredentials.class, name = "NO_SEC"), + @JsonSubTypes.Type(value = PSKServerCredentials.class, name = "PSK"), + @JsonSubTypes.Type(value = RPKServerCredentials.class, name = "RPK"), + @JsonSubTypes.Type(value = X509ServerCredentials.class, name = "X509") +}) +@JsonIgnoreProperties(ignoreUnknown = true) +public interface LwM2MServerCredentials { + + @JsonIgnore + LwM2MSecurityMode getSecurityMode(); +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/NoSecServerCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/NoSecServerCredentials.java new file mode 100644 index 0000000000..a7a76e5192 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/NoSecServerCredentials.java @@ -0,0 +1,24 @@ +/** + * 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.device.credentials.lwm2m; + +public class NoSecServerCredentials implements LwM2MServerCredentials { + + @Override + public LwM2MSecurityMode getSecurityMode() { + return LwM2MSecurityMode.NO_SEC; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKServerCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKServerCredentials.java new file mode 100644 index 0000000000..ab9e515482 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKServerCredentials.java @@ -0,0 +1,24 @@ +/** + * 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.device.credentials.lwm2m; + +public class PSKServerCredentials extends AbstractLwM2MServerCredentialsWithKeys { + + @Override + public LwM2MSecurityMode getSecurityMode() { + return LwM2MSecurityMode.PSK; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKServerCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKServerCredentials.java new file mode 100644 index 0000000000..f18e74f46c --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKServerCredentials.java @@ -0,0 +1,24 @@ +/** + * 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.device.credentials.lwm2m; + +public class RPKServerCredentials extends AbstractLwM2MServerCredentialsWithKeys { + + @Override + public LwM2MSecurityMode getSecurityMode() { + return LwM2MSecurityMode.RPK; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/X509ServerCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/X509ServerCredentials.java new file mode 100644 index 0000000000..a7b98608e2 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/X509ServerCredentials.java @@ -0,0 +1,24 @@ +/** + * 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.device.credentials.lwm2m; + +public class X509ServerCredentials extends AbstractLwM2MServerCredentialsWithKeys { + + @Override + public LwM2MSecurityMode getSecurityMode() { + return LwM2MSecurityMode.X509; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/Edge.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/Edge.java index 939bfb3648..7a858f8325 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/Edge.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/Edge.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.common.data.edge; -import com.fasterxml.jackson.databind.JsonNode; import lombok.EqualsAndHashCode; import lombok.Getter; import lombok.Setter; @@ -70,6 +69,19 @@ public class Edge extends SearchTextBasedWithAdditionalInfo implements H this.cloudEndpoint = edge.getCloudEndpoint(); } + public void update(Edge edge) { + this.tenantId = edge.getTenantId(); + this.customerId = edge.getCustomerId(); + this.rootRuleChainId = edge.getRootRuleChainId(); + this.type = edge.getType(); + this.label = edge.getLabel(); + this.name = edge.getName(); + this.routingKey = edge.getRoutingKey(); + this.secret = edge.getSecret(); + this.edgeLicenseKey = edge.getEdgeLicenseKey(); + this.cloudEndpoint = edge.getCloudEndpoint(); + } + @Override public String getSearchText() { return getName(); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java index 473429d526..be4143e388 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java @@ -577,6 +577,14 @@ public class JsonConverter { return GSON.toJson(element); } + public static JsonObject toJsonObject(Object o) { + return (JsonObject) GSON.toJsonTree(o); + } + + public static T fromJson(JsonElement element, Class type) { + return GSON.fromJson(element, type); + } + public static void setTypeCastEnabled(boolean enabled) { isTypeCastEnabled = enabled; } diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 18b7abb67e..da50963bd4 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -118,4 +118,9 @@ public class JacksonUtil { public static JsonNode valueToTree(T value) { return OBJECT_MAPPER.valueToTree(value); } + + public static T treeToValue(JsonNode tree, Class type) throws JsonProcessingException { + return OBJECT_MAPPER.treeToValue(tree, type); + } + } diff --git a/dao/pom.xml b/dao/pom.xml index c4282d82ab..922b389bff 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -227,6 +227,10 @@ org.elasticsearch.client rest + + org.eclipse.leshan + leshan-core + diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java index 3cce5e1506..4b9b4941e6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java @@ -16,20 +16,28 @@ package org.thingsboard.server.dao.device; -import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.codec.binary.Hex; +import org.eclipse.leshan.core.util.SecurityUtil; import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.annotation.CacheEvict; import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; -import org.springframework.util.StringUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MBootstrapCredentials; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MClientCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MServerCredentials; import org.thingsboard.server.common.data.device.credentials.lwm2m.PSKClientCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.PSKServerCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.RPKClientCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.RPKServerCredentials; import org.thingsboard.server.common.data.device.credentials.lwm2m.X509ClientCredentials; +import org.thingsboard.server.common.data.device.credentials.lwm2m.X509ServerCredentials; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -37,6 +45,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.msg.EncryptionUtil; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException; import org.thingsboard.server.dao.service.DataValidator; import static org.thingsboard.server.common.data.CacheConstants.DEVICE_CREDENTIALS_CACHE; @@ -83,17 +92,7 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen if (deviceCredentials.getCredentialsType() == null) { throw new DataValidationException("Device credentials type should be specified"); } - switch (deviceCredentials.getCredentialsType()) { - case X509_CERTIFICATE: - formatCertData(deviceCredentials); - break; - case MQTT_BASIC: - formatSimpleMqttCredentials(deviceCredentials); - break; - case LWM2M_CREDENTIALS: - formatSimpleLwm2mCredentials(deviceCredentials); - break; - } + formatCredentials(deviceCredentials); log.trace("Executing updateDeviceCredentials [{}]", deviceCredentials); credentialsValidator.validate(deviceCredentials, id -> tenantId); try { @@ -109,6 +108,21 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen } } + @Override + public void formatCredentials(DeviceCredentials deviceCredentials) { + switch (deviceCredentials.getCredentialsType()) { + case X509_CERTIFICATE: + formatCertData(deviceCredentials); + break; + case MQTT_BASIC: + formatSimpleMqttCredentials(deviceCredentials); + break; + case LWM2M_CREDENTIALS: + formatSimpleLwm2mCredentials(deviceCredentials); + break; + } + } + private void formatSimpleMqttCredentials(DeviceCredentials deviceCredentials) { BasicMqttCredentials mqttCredentials; try { @@ -117,11 +131,16 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen throw new IllegalArgumentException(); } } catch (IllegalArgumentException e) { - throw new DataValidationException("Invalid credentials body for simple mqtt credentials!"); + throw new DeviceCredentialsValidationException("Invalid credentials body for simple mqtt credentials!"); } + if (StringUtils.isEmpty(mqttCredentials.getClientId()) && StringUtils.isEmpty(mqttCredentials.getUserName())) { - throw new DataValidationException("Both mqtt client id and user name are empty!"); + throw new DeviceCredentialsValidationException("Both mqtt client id and user name are empty!"); } + if (StringUtils.isNotEmpty(mqttCredentials.getClientId()) && StringUtils.isNotEmpty(mqttCredentials.getPassword()) && StringUtils.isEmpty(mqttCredentials.getUserName())) { + throw new DeviceCredentialsValidationException("Password cannot be specified along with client id"); + } + if (StringUtils.isEmpty(mqttCredentials.getClientId())) { deviceCredentials.setCredentialsId(mqttCredentials.getUserName()); } else if (StringUtils.isEmpty(mqttCredentials.getUserName())) { @@ -129,7 +148,7 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen } else { deviceCredentials.setCredentialsId(EncryptionUtil.getSha3Hash("|", mqttCredentials.getClientId(), mqttCredentials.getUserName())); } - if (!StringUtils.isEmpty(mqttCredentials.getPassword())) { + if (StringUtils.isNotEmpty(mqttCredentials.getPassword())) { mqttCredentials.setPassword(mqttCredentials.getPassword()); } deviceCredentials.setCredentialsValue(JacksonUtil.toString(mqttCredentials)); @@ -143,22 +162,16 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen } private void formatSimpleLwm2mCredentials(DeviceCredentials deviceCredentials) { - LwM2MClientCredentials clientCredentials; - ObjectNode json; + LwM2MDeviceCredentials lwM2MCredentials; try { - json = JacksonUtil.fromString(deviceCredentials.getCredentialsValue(), ObjectNode.class); - if (json == null) { - throw new IllegalArgumentException(); - } - clientCredentials = JacksonUtil.convertValue(json.get("client"), LwM2MClientCredentials.class); - if (clientCredentials == null) { - throw new IllegalArgumentException(); - } + lwM2MCredentials = JacksonUtil.fromString(deviceCredentials.getCredentialsValue(), LwM2MDeviceCredentials.class); + validateLwM2MDeviceCredentials(lwM2MCredentials); } catch (IllegalArgumentException e) { - throw new DataValidationException("Invalid credentials body for LwM2M credentials!"); + throw new DeviceCredentialsValidationException("Invalid credentials body for LwM2M credentials!"); } String credentialsId = null; + LwM2MClientCredentials clientCredentials = lwM2MCredentials.getClient(); switch (clientCredentials.getSecurityConfigClientMode()) { case NO_SEC: @@ -174,8 +187,8 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen String cert = EncryptionUtil.trimNewLines(x509Config.getCert()); String sha3Hash = EncryptionUtil.getSha3Hash(cert); x509Config.setCert(cert); - ((ObjectNode) json.get("client")).put("cert", cert); - deviceCredentials.setCredentialsValue(JacksonUtil.toString(json)); + ((X509ClientCredentials) clientCredentials).setCert(cert); + deviceCredentials.setCredentialsValue(JacksonUtil.toString(lwM2MCredentials)); credentialsId = sha3Hash; } else { credentialsId = x509Config.getEndpoint(); @@ -183,11 +196,163 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen break; } if (credentialsId == null) { - throw new DataValidationException("Invalid credentials body for LwM2M credentials!"); + throw new DeviceCredentialsValidationException("Invalid credentials body for LwM2M credentials!"); } deviceCredentials.setCredentialsId(credentialsId); } + private void validateLwM2MDeviceCredentials(LwM2MDeviceCredentials lwM2MCredentials) { + if (lwM2MCredentials == null) { + throw new DeviceCredentialsValidationException("LwM2M credentials should be specified!"); + } + + LwM2MClientCredentials clientCredentials = lwM2MCredentials.getClient(); + if (clientCredentials == null) { + throw new DeviceCredentialsValidationException("LwM2M client credentials should be specified!"); + } + validateLwM2MClientCredentials(clientCredentials); + + LwM2MBootstrapCredentials bootstrapCredentials = lwM2MCredentials.getBootstrap(); + if (bootstrapCredentials == null) { + throw new DeviceCredentialsValidationException("LwM2M bootstrap credentials should be specified!"); + } + + LwM2MServerCredentials bootstrapServerCredentials = bootstrapCredentials.getBootstrapServer(); + if (bootstrapServerCredentials == null) { + throw new DeviceCredentialsValidationException("LwM2M bootstrap server credentials should be specified!"); + } + validateServerCredentials(bootstrapServerCredentials, "Bootstrap server"); + + LwM2MServerCredentials lwm2mServerCredentials = bootstrapCredentials.getLwm2mServer(); + if (lwm2mServerCredentials == null) { + throw new DeviceCredentialsValidationException("LwM2M lwm2m server credentials should be specified!"); + } + validateServerCredentials(lwm2mServerCredentials, "LwM2M server"); + } + + private void validateLwM2MClientCredentials(LwM2MClientCredentials clientCredentials) { + if (StringUtils.isEmpty(clientCredentials.getEndpoint())) { + throw new DeviceCredentialsValidationException("LwM2M client endpoint should be specified!"); + } + + switch (clientCredentials.getSecurityConfigClientMode()) { + case NO_SEC: + break; + case PSK: + PSKClientCredentials pskCredentials = (PSKClientCredentials) clientCredentials; + if (StringUtils.isEmpty(pskCredentials.getIdentity())) { + throw new DeviceCredentialsValidationException("LwM2M client PSK identity should be specified!"); + } + + String pskKey = pskCredentials.getKey(); + if (StringUtils.isEmpty(pskKey)) { + throw new DeviceCredentialsValidationException("LwM2M client PSK key should be specified!"); + } + + if (!pskKey.matches("-?[0-9a-fA-F]+")) { + throw new DeviceCredentialsValidationException("LwM2M client PSK key should be HexDecimal format!"); + } + + if (pskKey.length() % 32 != 0 || pskKey.length() > 128) { + throw new DeviceCredentialsValidationException("LwM2M client PSK key must be 32, 64, 128 characters!"); + } + break; + case RPK: + RPKClientCredentials rpkCredentials = (RPKClientCredentials) clientCredentials; + + if (StringUtils.isEmpty(rpkCredentials.getKey())) { + throw new DeviceCredentialsValidationException("LwM2M client RPK key should be specified!"); + } + + try { + SecurityUtil.publicKey.decode(rpkCredentials.getDecodedKey()); + } catch (Exception e) { + throw new DeviceCredentialsValidationException("LwM2M client RPK key should be in RFC7250 standard!"); + } + break; + case X509: + X509ClientCredentials x509CCredentials = (X509ClientCredentials) clientCredentials; + if (x509CCredentials.getCert() != null) { + try { + SecurityUtil.certificate.decode(Hex.decodeHex(x509CCredentials.getCert().toLowerCase().toCharArray())); + } catch (Exception e) { + throw new DeviceCredentialsValidationException("LwM2M client X509 certificate should be in DER-encoded X.509 format!"); + } + } + break; + } + } + + private void validateServerCredentials(LwM2MServerCredentials serverCredentials, String server) { + switch (serverCredentials.getSecurityMode()) { + case NO_SEC: + break; + case PSK: + PSKServerCredentials pskCredentials = (PSKServerCredentials) serverCredentials; + if (StringUtils.isEmpty(pskCredentials.getClientPublicKeyOrId())) { + throw new DeviceCredentialsValidationException(server + " client PSK public key or id should be specified!"); + } + + String pskKey = pskCredentials.getClientSecretKey(); + if (StringUtils.isEmpty(pskKey)) { + throw new DeviceCredentialsValidationException(server + " client PSK key should be specified!"); + } + + if (!pskKey.matches("-?[0-9a-fA-F]+")) { + throw new DeviceCredentialsValidationException(server + " client PSK key should be HexDecimal format!"); + } + + if (pskKey.length() % 32 != 0 || pskKey.length() > 128) { + throw new DeviceCredentialsValidationException(server + " client PSK key must be 32, 64, 128 characters!"); + } + break; + case RPK: + RPKServerCredentials rpkCredentials = (RPKServerCredentials) serverCredentials; + + if (StringUtils.isEmpty(rpkCredentials.getClientPublicKeyOrId())) { + throw new DeviceCredentialsValidationException(server + " client RPK public key or id should be specified!"); + } + + try { + SecurityUtil.publicKey.decode(rpkCredentials.getDecodedClientPublicKeyOrId()); + } catch (Exception e) { + throw new DeviceCredentialsValidationException(server + " client RPK public key or id should be in RFC7250 standard!"); + } + + if (StringUtils.isEmpty(rpkCredentials.getClientSecretKey())) { + throw new DeviceCredentialsValidationException(server + " client RPK secret key should be specified!"); + } + + try { + SecurityUtil.privateKey.decode(rpkCredentials.getDecodedClientSecretKey()); + } catch (Exception e) { + throw new DeviceCredentialsValidationException(server + " client RPK secret key should be in RFC5958 standard!"); + } + break; + case X509: + X509ServerCredentials x509CCredentials = (X509ServerCredentials) serverCredentials; + if (StringUtils.isEmpty(x509CCredentials.getClientPublicKeyOrId())) { + throw new DeviceCredentialsValidationException(server + " client X509 public key or id should be specified!"); + } + + try { + SecurityUtil.certificate.decode(x509CCredentials.getDecodedClientPublicKeyOrId()); + } catch (Exception e) { + throw new DeviceCredentialsValidationException(server + " client X509 public key or id should be in DER-encoded X.509 format!"); + } + if (StringUtils.isEmpty(x509CCredentials.getClientSecretKey())) { + throw new DeviceCredentialsValidationException(server + " client X509 secret key should be specified!"); + } + + try { + SecurityUtil.privateKey.decode(x509CCredentials.getDecodedClientSecretKey()); + } catch (Exception e) { + throw new DeviceCredentialsValidationException(server + " client X509 secret key should be in RFC5958 standard!"); + } + break; + } + } + @Override @CacheEvict(cacheNames = DEVICE_CREDENTIALS_CACHE, key = "'deviceCredentials_' + #deviceCredentials.credentialsId") public void deleteDeviceCredentials(TenantId tenantId, DeviceCredentials deviceCredentials) { @@ -201,38 +366,38 @@ public class DeviceCredentialsServiceImpl extends AbstractEntityService implemen @Override protected void validateCreate(TenantId tenantId, DeviceCredentials deviceCredentials) { if (deviceCredentialsDao.findByDeviceId(tenantId, deviceCredentials.getDeviceId().getId()) != null) { - throw new DataValidationException("Credentials for this device are already specified!"); + throw new DeviceCredentialsValidationException("Credentials for this device are already specified!"); } if (deviceCredentialsDao.findByCredentialsId(tenantId, deviceCredentials.getCredentialsId()) != null) { - throw new DataValidationException("Device credentials are already assigned to another device!"); + throw new DeviceCredentialsValidationException("Device credentials are already assigned to another device!"); } } @Override protected void validateUpdate(TenantId tenantId, DeviceCredentials deviceCredentials) { if (deviceCredentialsDao.findById(tenantId, deviceCredentials.getUuidId()) == null) { - throw new DataValidationException("Unable to update non-existent device credentials!"); + throw new DeviceCredentialsValidationException("Unable to update non-existent device credentials!"); } DeviceCredentials existingCredentials = deviceCredentialsDao.findByCredentialsId(tenantId, deviceCredentials.getCredentialsId()); if (existingCredentials != null && !existingCredentials.getId().equals(deviceCredentials.getId())) { - throw new DataValidationException("Device credentials are already assigned to another device!"); + throw new DeviceCredentialsValidationException("Device credentials are already assigned to another device!"); } } @Override protected void validateDataImpl(TenantId tenantId, DeviceCredentials deviceCredentials) { if (deviceCredentials.getDeviceId() == null) { - throw new DataValidationException("Device credentials should be assigned to device!"); + throw new DeviceCredentialsValidationException("Device credentials should be assigned to device!"); } if (deviceCredentials.getCredentialsType() == null) { - throw new DataValidationException("Device credentials type should be specified!"); + throw new DeviceCredentialsValidationException("Device credentials type should be specified!"); } if (StringUtils.isEmpty(deviceCredentials.getCredentialsId())) { - throw new DataValidationException("Device credentials id should be specified!"); + throw new DeviceCredentialsValidationException("Device credentials id should be specified!"); } Device device = deviceService.findDeviceById(tenantId, deviceCredentials.getDeviceId()); if (device == null) { - throw new DataValidationException("Can't assign device credentials to non-existent device!"); + throw new DeviceCredentialsValidationException("Can't assign device credentials to non-existent device!"); } } }; diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 69ccb2dfde..3255bf75b4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -228,6 +228,7 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe if (foundDeviceCredentials == null) { deviceCredentialsService.createDeviceCredentials(savedDevice.getTenantId(), deviceCredentials); } else { + deviceCredentials.setId(foundDeviceCredentials.getId()); deviceCredentialsService.updateDeviceCredentials(device.getTenantId(), deviceCredentials); } } @@ -241,7 +242,7 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe deviceCredentials.setDeviceId(new DeviceId(savedDevice.getUuidId())); deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); deviceCredentials.setCredentialsId(!StringUtils.isEmpty(accessToken) ? accessToken : RandomStringUtils.randomAlphanumeric(20)); - deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials); + deviceCredentialsService.createDeviceCredentials(savedDevice.getTenantId(), deviceCredentials); } return savedDevice; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/exception/DeviceCredentialsValidationException.java b/dao/src/main/java/org/thingsboard/server/dao/exception/DeviceCredentialsValidationException.java new file mode 100644 index 0000000000..d4f0b584eb --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/exception/DeviceCredentialsValidationException.java @@ -0,0 +1,22 @@ +/** + * 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.exception; + +public class DeviceCredentialsValidationException extends DataValidationException { + public DeviceCredentialsValidationException(String message) { + super(message); + } +} diff --git a/ui-ngx/src/app/core/http/asset.service.ts b/ui-ngx/src/app/core/http/asset.service.ts index aa4bb161d1..e106600c57 100644 --- a/ui-ngx/src/app/core/http/asset.service.ts +++ b/ui-ngx/src/app/core/http/asset.service.ts @@ -22,6 +22,7 @@ import { PageLink } from '@shared/models/page/page-link'; import { PageData } from '@shared/models/page/page-data'; import { EntitySubtype } from '@app/shared/models/entity-type.models'; import { Asset, AssetInfo, AssetSearchQuery } from '@app/shared/models/asset.models'; +import { BulkImportRequest, BulkImportResult } from '@home/components/import-export/import-export.models'; @Injectable({ providedIn: 'root' @@ -105,4 +106,8 @@ export class AssetService { defaultHttpOptionsFromConfig(config)); } + public bulkImportAssets(entitiesData: BulkImportRequest, config?: RequestConfig): Observable { + return this.http.post('/api/asset/bulk_import', entitiesData, defaultHttpOptionsFromConfig(config)); + } + } diff --git a/ui-ngx/src/app/core/http/device.service.ts b/ui-ngx/src/app/core/http/device.service.ts index c881e038e9..c8ea7e65ff 100644 --- a/ui-ngx/src/app/core/http/device.service.ts +++ b/ui-ngx/src/app/core/http/device.service.ts @@ -30,6 +30,7 @@ import { } from '@app/shared/models/device.models'; import { EntitySubtype } from '@app/shared/models/entity-type.models'; import { AuthService } from '@core/auth/auth.service'; +import { BulkImportRequest, BulkImportResult } from '@home/components/import-export/import-export.models'; import { PersistentRpc } from '@shared/models/rpc.models'; @Injectable({ @@ -178,4 +179,8 @@ export class DeviceService { defaultHttpOptionsFromConfig(config)); } + public bulkImportDevices(entitiesData: BulkImportRequest, config?: RequestConfig): Observable { + return this.http.post('/api/device/bulk_import', entitiesData, defaultHttpOptionsFromConfig(config)); + } + } diff --git a/ui-ngx/src/app/core/http/edge.service.ts b/ui-ngx/src/app/core/http/edge.service.ts index 7e4974aec4..9bec9b873c 100644 --- a/ui-ngx/src/app/core/http/edge.service.ts +++ b/ui-ngx/src/app/core/http/edge.service.ts @@ -23,6 +23,7 @@ import { PageData } from '@shared/models/page/page-data'; import { EntitySubtype } from '@app/shared/models/entity-type.models'; import { Edge, EdgeEvent, EdgeInfo, EdgeSearchQuery } from '@shared/models/edge.models'; import { EntityId } from '@shared/models/id/entity-id'; +import { BulkImportRequest, BulkImportResult } from '@home/components/import-export/import-export.models'; @Injectable({ providedIn: 'root' @@ -59,7 +60,7 @@ export class EdgeService { } public getCustomerEdgeInfos(customerId: string, pageLink: PageLink, type: string = '', - config?: RequestConfig): Observable> { + config?: RequestConfig): Observable> { return this.http.get>(`/api/customer/${customerId}/edgeInfos${pageLink.toQuery()}&type=${type}`, defaultHttpOptionsFromConfig(config)); } @@ -108,4 +109,8 @@ export class EdgeService { public findByName(edgeName: string, config?: RequestConfig): Observable { return this.http.get(`/api/tenant/edges?edgeName=${edgeName}`, defaultHttpOptionsFromConfig(config)); } + + public bulkImportEdges(entitiesData: BulkImportRequest, config?: RequestConfig): Observable { + return this.http.post('/api/edge/bulk_import', entitiesData, defaultHttpOptionsFromConfig(config)); + } } diff --git a/ui-ngx/src/app/core/http/entity.service.ts b/ui-ngx/src/app/core/http/entity.service.ts index 78a7dc76df..6665beb3b3 100644 --- a/ui-ngx/src/app/core/http/entity.service.ts +++ b/ui-ngx/src/app/core/http/entity.service.ts @@ -59,7 +59,7 @@ import { ImportEntityData } from '@shared/models/entity.models'; import { EntityRelationService } from '@core/http/entity-relation.service'; -import { deepClone, generateSecret, guid, isDefined, isDefinedAndNotNull } from '@core/utils'; +import { deepClone, generateSecret, guid, isDefined, isDefinedAndNotNull, isNotEmptyStr } from '@core/utils'; import { Asset } from '@shared/models/asset.models'; import { Device, DeviceCredentialsType } from '@shared/models/device.models'; import { AttributeService } from '@core/http/attribute.service'; @@ -964,7 +964,12 @@ export class EntityService { map(() => { return { create: { entity: 1 } } as ImportEntitiesResultInfo; }), - catchError(err => of({ error: { entity: 1 } } as ImportEntitiesResultInfo)) + catchError(err => of({ + error: { + entity: 1, + errors: err.message + } + } as ImportEntitiesResultInfo)) ); }), catchError(err => { @@ -988,13 +993,28 @@ export class EntityService { map(() => { return { update: { entity: 1 } } as ImportEntitiesResultInfo; }), - catchError(updateError => of({ error: { entity: 1 } } as ImportEntitiesResultInfo)) + catchError(updateError => of({ + error: { + entity: 1, + errors: updateError.message + } + } as ImportEntitiesResultInfo)) ); }), - catchError(findErr => of({ error: { entity: 1 } } as ImportEntitiesResultInfo)) + catchError(findErr => of({ + error: { + entity: 1, + errors: `Line: ${entityData.lineNumber}; Error: ${findErr.error.message}` + } + } as ImportEntitiesResultInfo)) ); } else { - return of({ error: { entity: 1 } } as ImportEntitiesResultInfo); + return of({ + error: { + entity: 1, + errors: `Line: ${entityData.lineNumber}; Error: ${err.error.message}` + } + } as ImportEntitiesResultInfo); } }) ); @@ -1050,7 +1070,6 @@ export class EntityService { break; } return saveEntityObservable; - } private getUpdateEntityTasks(entityType: EntityType, entityData: ImportEntityData | EdgeImportEntityData, @@ -1123,15 +1142,31 @@ export class EntityService { public saveEntityData(entityId: EntityId, entityData: ImportEntityData, config?: RequestConfig): Observable { const observables: Observable[] = []; let observable: Observable; - if (entityData.accessToken && entityData.accessToken !== '') { + if (Object.keys(entityData.credential).length) { + let credentialsType: DeviceCredentialsType; + let credentialsId: string = null; + let credentialsValue: string = null; + if (isDefinedAndNotNull(entityData.credential.mqtt)) { + credentialsType = DeviceCredentialsType.MQTT_BASIC; + credentialsValue = JSON.stringify(entityData.credential.mqtt); + } else if (isDefinedAndNotNull(entityData.credential.lwm2m)) { + credentialsType = DeviceCredentialsType.LWM2M_CREDENTIALS; + credentialsValue = JSON.stringify(entityData.credential.lwm2m); + } else if (isNotEmptyStr(entityData.credential.x509)) { + credentialsType = DeviceCredentialsType.X509_CERTIFICATE; + credentialsValue = entityData.credential.x509; + } else { + credentialsType = DeviceCredentialsType.ACCESS_TOKEN; + credentialsId = entityData.credential.accessToken; + } observable = this.deviceService.getDeviceCredentials(entityId.id, false, config).pipe( mergeMap((credentials) => { - credentials.credentialsId = entityData.accessToken; - credentials.credentialsType = DeviceCredentialsType.ACCESS_TOKEN; - credentials.credentialsValue = null; + credentials.credentialsId = credentialsId; + credentials.credentialsType = credentialsType; + credentials.credentialsValue = credentialsValue; return this.deviceService.saveDeviceCredentials(credentials, config).pipe( map(() => 'ok'), - catchError(err => of('error')) + catchError(err => of(`Line: ${entityData.lineNumber}; Error: ${err.error.message}`)) ); }) ); @@ -1141,7 +1176,7 @@ export class EntityService { observable = this.attributeService.saveEntityAttributes(entityId, AttributeScope.SHARED_SCOPE, entityData.attributes.shared, config).pipe( map(() => 'ok'), - catchError(err => of('error')) + catchError(err => of(`Line: ${entityData.lineNumber}; Error: ${err.error.message}`)) ); observables.push(observable); } @@ -1149,23 +1184,23 @@ export class EntityService { observable = this.attributeService.saveEntityAttributes(entityId, AttributeScope.SERVER_SCOPE, entityData.attributes.server, config).pipe( map(() => 'ok'), - catchError(err => of('error')) + catchError(err => of(`Line: ${entityData.lineNumber}; Error: ${err.error.message}`)) ); observables.push(observable); } if (entityData.timeseries && entityData.timeseries.length) { observable = this.attributeService.saveEntityTimeseries(entityId, 'time', entityData.timeseries, config).pipe( map(() => 'ok'), - catchError(err => of('error')) + catchError(err => of(`Line: ${entityData.lineNumber}; Error: ${err.error.message}`)) ); observables.push(observable); } if (observables.length) { return forkJoin(observables).pipe( map((response) => { - const hasError = response.filter((status) => status === 'error').length > 0; - if (hasError) { - throw Error(); + const hasError = response.filter((status) => status !== 'ok'); + if (hasError.length > 0) { + throw Error(hasError.join('\n')); } else { return response; } diff --git a/ui-ngx/src/app/modules/home/components/import-export/import-dialog-csv.component.html b/ui-ngx/src/app/modules/home/components/import-export/import-dialog-csv.component.html index 77f9a43e66..035e43690c 100644 --- a/ui-ngx/src/app/modules/home/components/import-export/import-dialog-csv.component.html +++ b/ui-ngx/src/app/modules/home/components/import-export/import-dialog-csv.component.html @@ -94,7 +94,7 @@
{{ 'import.stepper-text.column-type' | translate }} - +