1000 changed files with 53784 additions and 13752 deletions
File diff suppressed because it is too large
File diff suppressed because it is too large
@ -0,0 +1,24 @@ |
|||||
|
{ |
||||
|
"providerId": "Apple", |
||||
|
"additionalInfo": null, |
||||
|
"accessTokenUri": "https://appleid.apple.com/auth/token", |
||||
|
"authorizationUri": "https://appleid.apple.com/auth/authorize?response_mode=form_post", |
||||
|
"scope": ["email","openid","name"], |
||||
|
"jwkSetUri": "https://appleid.apple.com/auth/keys", |
||||
|
"userInfoUri": null, |
||||
|
"clientAuthenticationMethod": "POST", |
||||
|
"userNameAttributeName": "email", |
||||
|
"mapperConfig": { |
||||
|
"type": "APPLE", |
||||
|
"basic": { |
||||
|
"emailAttributeKey": "email", |
||||
|
"firstNameAttributeKey": "firstName", |
||||
|
"lastNameAttributeKey": "lastName", |
||||
|
"tenantNameStrategy": "DOMAIN" |
||||
|
} |
||||
|
}, |
||||
|
"comment": null, |
||||
|
"loginButtonIcon": "apple-logo", |
||||
|
"loginButtonLabel": "Apple", |
||||
|
"helpLink": "https://developer.apple.com/sign-in-with-apple/get-started/" |
||||
|
} |
||||
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@ -0,0 +1,225 @@ |
|||||
|
/** |
||||
|
* 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.controller; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.apache.commons.lang3.StringUtils; |
||||
|
import org.springframework.core.io.ByteArrayResource; |
||||
|
import org.springframework.http.HttpHeaders; |
||||
|
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.RequestBody; |
||||
|
import org.springframework.web.bind.annotation.RequestMapping; |
||||
|
import org.springframework.web.bind.annotation.RequestMethod; |
||||
|
import org.springframework.web.bind.annotation.RequestParam; |
||||
|
import org.springframework.web.bind.annotation.ResponseBody; |
||||
|
import org.springframework.web.bind.annotation.RestController; |
||||
|
import org.springframework.web.multipart.MultipartFile; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
import org.thingsboard.server.common.data.OtaPackage; |
||||
|
import org.thingsboard.server.common.data.OtaPackageInfo; |
||||
|
import org.thingsboard.server.common.data.audit.ActionType; |
||||
|
import org.thingsboard.server.common.data.exception.ThingsboardException; |
||||
|
import org.thingsboard.server.common.data.ota.ChecksumAlgorithm; |
||||
|
import org.thingsboard.server.common.data.ota.OtaPackageType; |
||||
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
||||
|
import org.thingsboard.server.common.data.id.OtaPackageId; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.security.permission.Operation; |
||||
|
import org.thingsboard.server.service.security.permission.Resource; |
||||
|
|
||||
|
import java.nio.ByteBuffer; |
||||
|
|
||||
|
@Slf4j |
||||
|
@RestController |
||||
|
@TbCoreComponent |
||||
|
@RequestMapping("/api") |
||||
|
public class OtaPackageController extends BaseController { |
||||
|
|
||||
|
public static final String OTA_PACKAGE_ID = "otaPackageId"; |
||||
|
public static final String CHECKSUM_ALGORITHM = "checksumAlgorithm"; |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority( 'TENANT_ADMIN')") |
||||
|
@RequestMapping(value = "/otaPackage/{otaPackageId}/download", method = RequestMethod.GET) |
||||
|
@ResponseBody |
||||
|
public ResponseEntity<org.springframework.core.io.Resource> downloadOtaPackage(@PathVariable(OTA_PACKAGE_ID) String strOtaPackageId) throws ThingsboardException { |
||||
|
checkParameter(OTA_PACKAGE_ID, strOtaPackageId); |
||||
|
try { |
||||
|
OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); |
||||
|
OtaPackage otaPackage = checkOtaPackageId(otaPackageId, Operation.READ); |
||||
|
|
||||
|
if (otaPackage.hasUrl()) { |
||||
|
return ResponseEntity.badRequest().build(); |
||||
|
} |
||||
|
|
||||
|
ByteArrayResource resource = new ByteArrayResource(otaPackage.getData().array()); |
||||
|
return ResponseEntity.ok() |
||||
|
.header(HttpHeaders.CONTENT_DISPOSITION, "attachment;filename=" + otaPackage.getFileName()) |
||||
|
.header("x-filename", otaPackage.getFileName()) |
||||
|
.contentLength(resource.contentLength()) |
||||
|
.contentType(parseMediaType(otaPackage.getContentType())) |
||||
|
.body(resource); |
||||
|
} catch (Exception e) { |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
||||
|
@RequestMapping(value = "/otaPackage/info/{otaPackageId}", method = RequestMethod.GET) |
||||
|
@ResponseBody |
||||
|
public OtaPackageInfo getOtaPackageInfoById(@PathVariable(OTA_PACKAGE_ID) String strOtaPackageId) throws ThingsboardException { |
||||
|
checkParameter(OTA_PACKAGE_ID, strOtaPackageId); |
||||
|
try { |
||||
|
OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); |
||||
|
return checkNotNull(otaPackageService.findOtaPackageInfoById(getTenantId(), otaPackageId)); |
||||
|
} catch (Exception e) { |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
||||
|
@RequestMapping(value = "/otaPackage/{otaPackageId}", method = RequestMethod.GET) |
||||
|
@ResponseBody |
||||
|
public OtaPackage getOtaPackageById(@PathVariable(OTA_PACKAGE_ID) String strOtaPackageId) throws ThingsboardException { |
||||
|
checkParameter(OTA_PACKAGE_ID, strOtaPackageId); |
||||
|
try { |
||||
|
OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); |
||||
|
return checkOtaPackageId(otaPackageId, Operation.READ); |
||||
|
} catch (Exception e) { |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
||||
|
@RequestMapping(value = "/otaPackage", method = RequestMethod.POST) |
||||
|
@ResponseBody |
||||
|
public OtaPackageInfo saveOtaPackageInfo(@RequestBody OtaPackageInfo otaPackageInfo) throws ThingsboardException { |
||||
|
boolean created = otaPackageInfo.getId() == null; |
||||
|
try { |
||||
|
otaPackageInfo.setTenantId(getTenantId()); |
||||
|
checkEntity(otaPackageInfo.getId(), otaPackageInfo, Resource.OTA_PACKAGE); |
||||
|
OtaPackageInfo savedOtaPackageInfo = otaPackageService.saveOtaPackageInfo(otaPackageInfo); |
||||
|
logEntityAction(savedOtaPackageInfo.getId(), savedOtaPackageInfo, |
||||
|
null, created ? ActionType.ADDED : ActionType.UPDATED, null); |
||||
|
return savedOtaPackageInfo; |
||||
|
} catch (Exception e) { |
||||
|
logEntityAction(emptyId(EntityType.OTA_PACKAGE), otaPackageInfo, |
||||
|
null, created ? ActionType.ADDED : ActionType.UPDATED, e); |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
||||
|
@RequestMapping(value = "/otaPackage/{otaPackageId}", method = RequestMethod.POST) |
||||
|
@ResponseBody |
||||
|
public OtaPackageInfo saveOtaPackageData(@PathVariable(OTA_PACKAGE_ID) String strOtaPackageId, |
||||
|
@RequestParam(required = false) String checksum, |
||||
|
@RequestParam(CHECKSUM_ALGORITHM) String checksumAlgorithmStr, |
||||
|
@RequestBody MultipartFile file) throws ThingsboardException { |
||||
|
checkParameter(OTA_PACKAGE_ID, strOtaPackageId); |
||||
|
checkParameter(CHECKSUM_ALGORITHM, checksumAlgorithmStr); |
||||
|
try { |
||||
|
OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); |
||||
|
OtaPackageInfo info = checkOtaPackageInfoId(otaPackageId, Operation.READ); |
||||
|
|
||||
|
OtaPackage otaPackage = new OtaPackage(otaPackageId); |
||||
|
otaPackage.setCreatedTime(info.getCreatedTime()); |
||||
|
otaPackage.setTenantId(getTenantId()); |
||||
|
otaPackage.setDeviceProfileId(info.getDeviceProfileId()); |
||||
|
otaPackage.setType(info.getType()); |
||||
|
otaPackage.setTitle(info.getTitle()); |
||||
|
otaPackage.setVersion(info.getVersion()); |
||||
|
otaPackage.setAdditionalInfo(info.getAdditionalInfo()); |
||||
|
|
||||
|
ChecksumAlgorithm checksumAlgorithm = ChecksumAlgorithm.valueOf(checksumAlgorithmStr.toUpperCase()); |
||||
|
|
||||
|
byte[] bytes = file.getBytes(); |
||||
|
if (StringUtils.isEmpty(checksum)) { |
||||
|
checksum = otaPackageService.generateChecksum(checksumAlgorithm, ByteBuffer.wrap(bytes)); |
||||
|
} |
||||
|
|
||||
|
otaPackage.setChecksumAlgorithm(checksumAlgorithm); |
||||
|
otaPackage.setChecksum(checksum); |
||||
|
otaPackage.setFileName(file.getOriginalFilename()); |
||||
|
otaPackage.setContentType(file.getContentType()); |
||||
|
otaPackage.setData(ByteBuffer.wrap(bytes)); |
||||
|
otaPackage.setDataSize((long) bytes.length); |
||||
|
OtaPackageInfo savedOtaPackage = otaPackageService.saveOtaPackage(otaPackage); |
||||
|
logEntityAction(savedOtaPackage.getId(), savedOtaPackage, null, ActionType.UPDATED, null); |
||||
|
return savedOtaPackage; |
||||
|
} catch (Exception e) { |
||||
|
logEntityAction(emptyId(EntityType.OTA_PACKAGE), null, null, ActionType.UPDATED, e, strOtaPackageId); |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
||||
|
@RequestMapping(value = "/otaPackages", method = RequestMethod.GET) |
||||
|
@ResponseBody |
||||
|
public PageData<OtaPackageInfo> getOtaPackages(@RequestParam int pageSize, |
||||
|
@RequestParam int page, |
||||
|
@RequestParam(required = false) String textSearch, |
||||
|
@RequestParam(required = false) String sortProperty, |
||||
|
@RequestParam(required = false) String sortOrder) throws ThingsboardException { |
||||
|
try { |
||||
|
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); |
||||
|
return checkNotNull(otaPackageService.findTenantOtaPackagesByTenantId(getTenantId(), pageLink)); |
||||
|
} catch (Exception e) { |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
||||
|
@RequestMapping(value = "/otaPackages/{deviceProfileId}/{type}", method = RequestMethod.GET) |
||||
|
@ResponseBody |
||||
|
public PageData<OtaPackageInfo> getOtaPackages(@PathVariable("deviceProfileId") String strDeviceProfileId, |
||||
|
@PathVariable("type") String strType, |
||||
|
@RequestParam int pageSize, |
||||
|
@RequestParam int page, |
||||
|
@RequestParam(required = false) String textSearch, |
||||
|
@RequestParam(required = false) String sortProperty, |
||||
|
@RequestParam(required = false) String sortOrder) throws ThingsboardException { |
||||
|
checkParameter("deviceProfileId", strDeviceProfileId); |
||||
|
checkParameter("type", strType); |
||||
|
try { |
||||
|
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); |
||||
|
return checkNotNull(otaPackageService.findTenantOtaPackagesByTenantIdAndDeviceProfileIdAndTypeAndHasData(getTenantId(), |
||||
|
new DeviceProfileId(toUUID(strDeviceProfileId)), OtaPackageType.valueOf(strType), pageLink)); |
||||
|
} catch (Exception e) { |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
||||
|
@RequestMapping(value = "/otaPackage/{otaPackageId}", method = RequestMethod.DELETE) |
||||
|
@ResponseBody |
||||
|
public void deleteOtaPackage(@PathVariable("otaPackageId") String strOtaPackageId) throws ThingsboardException { |
||||
|
checkParameter(OTA_PACKAGE_ID, strOtaPackageId); |
||||
|
try { |
||||
|
OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); |
||||
|
OtaPackageInfo info = checkOtaPackageInfoId(otaPackageId, Operation.DELETE); |
||||
|
otaPackageService.deleteOtaPackage(getTenantId(), otaPackageId); |
||||
|
logEntityAction(otaPackageId, info, null, ActionType.DELETED, null, strOtaPackageId); |
||||
|
} catch (Exception e) { |
||||
|
logEntityAction(emptyId(EntityType.OTA_PACKAGE), null, null, ActionType.DELETED, e, strOtaPackageId); |
||||
|
throw handleException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,256 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.action; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.ObjectMapper; |
||||
|
import com.fasterxml.jackson.databind.node.ArrayNode; |
||||
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.apache.commons.lang3.StringUtils; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.DataConstants; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
import org.thingsboard.server.common.data.HasName; |
||||
|
import org.thingsboard.server.common.data.HasTenantId; |
||||
|
import org.thingsboard.server.common.data.User; |
||||
|
import org.thingsboard.server.common.data.audit.ActionType; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.DataType; |
||||
|
import org.thingsboard.server.common.data.kv.KvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgDataType; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.queue.TbClusterService; |
||||
|
|
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
@TbCoreComponent |
||||
|
@Service |
||||
|
@RequiredArgsConstructor |
||||
|
@Slf4j |
||||
|
public class RuleEngineEntityActionService { |
||||
|
private final TbClusterService tbClusterService; |
||||
|
|
||||
|
private static final ObjectMapper json = new ObjectMapper(); |
||||
|
|
||||
|
public void pushEntityActionToRuleEngine(EntityId entityId, HasName entity, TenantId tenantId, CustomerId customerId, |
||||
|
ActionType actionType, User user, Object... additionalInfo) { |
||||
|
String msgType = null; |
||||
|
switch (actionType) { |
||||
|
case ADDED: |
||||
|
msgType = DataConstants.ENTITY_CREATED; |
||||
|
break; |
||||
|
case DELETED: |
||||
|
msgType = DataConstants.ENTITY_DELETED; |
||||
|
break; |
||||
|
case UPDATED: |
||||
|
msgType = DataConstants.ENTITY_UPDATED; |
||||
|
break; |
||||
|
case ASSIGNED_TO_CUSTOMER: |
||||
|
msgType = DataConstants.ENTITY_ASSIGNED; |
||||
|
break; |
||||
|
case UNASSIGNED_FROM_CUSTOMER: |
||||
|
msgType = DataConstants.ENTITY_UNASSIGNED; |
||||
|
break; |
||||
|
case ATTRIBUTES_UPDATED: |
||||
|
msgType = DataConstants.ATTRIBUTES_UPDATED; |
||||
|
break; |
||||
|
case ATTRIBUTES_DELETED: |
||||
|
msgType = DataConstants.ATTRIBUTES_DELETED; |
||||
|
break; |
||||
|
case ALARM_ACK: |
||||
|
msgType = DataConstants.ALARM_ACK; |
||||
|
break; |
||||
|
case ALARM_CLEAR: |
||||
|
msgType = DataConstants.ALARM_CLEAR; |
||||
|
break; |
||||
|
case ALARM_DELETE: |
||||
|
msgType = DataConstants.ALARM_DELETE; |
||||
|
break; |
||||
|
case ASSIGNED_FROM_TENANT: |
||||
|
msgType = DataConstants.ENTITY_ASSIGNED_FROM_TENANT; |
||||
|
break; |
||||
|
case ASSIGNED_TO_TENANT: |
||||
|
msgType = DataConstants.ENTITY_ASSIGNED_TO_TENANT; |
||||
|
break; |
||||
|
case PROVISION_SUCCESS: |
||||
|
msgType = DataConstants.PROVISION_SUCCESS; |
||||
|
break; |
||||
|
case PROVISION_FAILURE: |
||||
|
msgType = DataConstants.PROVISION_FAILURE; |
||||
|
break; |
||||
|
case TIMESERIES_UPDATED: |
||||
|
msgType = DataConstants.TIMESERIES_UPDATED; |
||||
|
break; |
||||
|
case TIMESERIES_DELETED: |
||||
|
msgType = DataConstants.TIMESERIES_DELETED; |
||||
|
break; |
||||
|
case ASSIGNED_TO_EDGE: |
||||
|
msgType = DataConstants.ENTITY_ASSIGNED_TO_EDGE; |
||||
|
break; |
||||
|
case UNASSIGNED_FROM_EDGE: |
||||
|
msgType = DataConstants.ENTITY_UNASSIGNED_FROM_EDGE; |
||||
|
break; |
||||
|
} |
||||
|
if (!StringUtils.isEmpty(msgType)) { |
||||
|
try { |
||||
|
TbMsgMetaData metaData = new TbMsgMetaData(); |
||||
|
if (user != null) { |
||||
|
metaData.putValue("userId", user.getId().toString()); |
||||
|
metaData.putValue("userName", user.getName()); |
||||
|
} |
||||
|
if (customerId != null && !customerId.isNullUid()) { |
||||
|
metaData.putValue("customerId", customerId.toString()); |
||||
|
} |
||||
|
if (actionType == ActionType.ASSIGNED_TO_CUSTOMER) { |
||||
|
String strCustomerId = extractParameter(String.class, 1, additionalInfo); |
||||
|
String strCustomerName = extractParameter(String.class, 2, additionalInfo); |
||||
|
metaData.putValue("assignedCustomerId", strCustomerId); |
||||
|
metaData.putValue("assignedCustomerName", strCustomerName); |
||||
|
} else if (actionType == ActionType.UNASSIGNED_FROM_CUSTOMER) { |
||||
|
String strCustomerId = extractParameter(String.class, 1, additionalInfo); |
||||
|
String strCustomerName = extractParameter(String.class, 2, additionalInfo); |
||||
|
metaData.putValue("unassignedCustomerId", strCustomerId); |
||||
|
metaData.putValue("unassignedCustomerName", strCustomerName); |
||||
|
} else if (actionType == ActionType.ASSIGNED_FROM_TENANT) { |
||||
|
String strTenantId = extractParameter(String.class, 0, additionalInfo); |
||||
|
String strTenantName = extractParameter(String.class, 1, additionalInfo); |
||||
|
metaData.putValue("assignedFromTenantId", strTenantId); |
||||
|
metaData.putValue("assignedFromTenantName", strTenantName); |
||||
|
} else if (actionType == ActionType.ASSIGNED_TO_TENANT) { |
||||
|
String strTenantId = extractParameter(String.class, 0, additionalInfo); |
||||
|
String strTenantName = extractParameter(String.class, 1, additionalInfo); |
||||
|
metaData.putValue("assignedToTenantId", strTenantId); |
||||
|
metaData.putValue("assignedToTenantName", strTenantName); |
||||
|
} else if (actionType == ActionType.ASSIGNED_TO_EDGE) { |
||||
|
String strEdgeId = extractParameter(String.class, 1, additionalInfo); |
||||
|
String strEdgeName = extractParameter(String.class, 2, additionalInfo); |
||||
|
metaData.putValue("assignedEdgeId", strEdgeId); |
||||
|
metaData.putValue("assignedEdgeName", strEdgeName); |
||||
|
} else if (actionType == ActionType.UNASSIGNED_FROM_EDGE) { |
||||
|
String strEdgeId = extractParameter(String.class, 1, additionalInfo); |
||||
|
String strEdgeName = extractParameter(String.class, 2, additionalInfo); |
||||
|
metaData.putValue("unassignedEdgeId", strEdgeId); |
||||
|
metaData.putValue("unassignedEdgeName", strEdgeName); |
||||
|
} |
||||
|
ObjectNode entityNode; |
||||
|
if (entity != null) { |
||||
|
entityNode = json.valueToTree(entity); |
||||
|
if (entityId.getEntityType() == EntityType.DASHBOARD) { |
||||
|
entityNode.put("configuration", ""); |
||||
|
} |
||||
|
} else { |
||||
|
entityNode = json.createObjectNode(); |
||||
|
if (actionType == ActionType.ATTRIBUTES_UPDATED) { |
||||
|
String scope = extractParameter(String.class, 0, additionalInfo); |
||||
|
@SuppressWarnings("unchecked") |
||||
|
List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo); |
||||
|
metaData.putValue(DataConstants.SCOPE, scope); |
||||
|
if (attributes != null) { |
||||
|
for (AttributeKvEntry attr : attributes) { |
||||
|
addKvEntry(entityNode, attr); |
||||
|
} |
||||
|
} |
||||
|
} else if (actionType == ActionType.ATTRIBUTES_DELETED) { |
||||
|
String scope = extractParameter(String.class, 0, additionalInfo); |
||||
|
@SuppressWarnings("unchecked") |
||||
|
List<String> keys = extractParameter(List.class, 1, additionalInfo); |
||||
|
metaData.putValue(DataConstants.SCOPE, scope); |
||||
|
ArrayNode attrsArrayNode = entityNode.putArray("attributes"); |
||||
|
if (keys != null) { |
||||
|
keys.forEach(attrsArrayNode::add); |
||||
|
} |
||||
|
} else if (actionType == ActionType.TIMESERIES_UPDATED) { |
||||
|
@SuppressWarnings("unchecked") |
||||
|
List<TsKvEntry> timeseries = extractParameter(List.class, 0, additionalInfo); |
||||
|
addTimeseries(entityNode, timeseries); |
||||
|
} else if (actionType == ActionType.TIMESERIES_DELETED) { |
||||
|
@SuppressWarnings("unchecked") |
||||
|
List<String> keys = extractParameter(List.class, 0, additionalInfo); |
||||
|
if (keys != null) { |
||||
|
ArrayNode timeseriesArrayNode = entityNode.putArray("timeseries"); |
||||
|
keys.forEach(timeseriesArrayNode::add); |
||||
|
} |
||||
|
entityNode.put("startTs", extractParameter(Long.class, 1, additionalInfo)); |
||||
|
entityNode.put("endTs", extractParameter(Long.class, 2, additionalInfo)); |
||||
|
} |
||||
|
} |
||||
|
TbMsg tbMsg = TbMsg.newMsg(msgType, entityId, customerId, metaData, TbMsgDataType.JSON, json.writeValueAsString(entityNode)); |
||||
|
if (tenantId.isNullUid()) { |
||||
|
if (entity instanceof HasTenantId) { |
||||
|
tenantId = ((HasTenantId) entity).getTenantId(); |
||||
|
} |
||||
|
} |
||||
|
tbClusterService.pushMsgToRuleEngine(tenantId, entityId, tbMsg, null); |
||||
|
} catch (Exception e) { |
||||
|
log.warn("[{}] Failed to push entity action to rule engine: {}", entityId, actionType, e); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
|
||||
|
private <T> T extractParameter(Class<T> clazz, int index, Object... additionalInfo) { |
||||
|
T result = null; |
||||
|
if (additionalInfo != null && additionalInfo.length > index) { |
||||
|
Object paramObject = additionalInfo[index]; |
||||
|
if (clazz.isInstance(paramObject)) { |
||||
|
result = clazz.cast(paramObject); |
||||
|
} |
||||
|
} |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
private void addTimeseries(ObjectNode entityNode, List<TsKvEntry> timeseries) throws Exception { |
||||
|
if (timeseries != null && !timeseries.isEmpty()) { |
||||
|
ArrayNode result = entityNode.putArray("timeseries"); |
||||
|
Map<Long, List<TsKvEntry>> groupedTelemetry = timeseries.stream() |
||||
|
.collect(Collectors.groupingBy(TsKvEntry::getTs)); |
||||
|
for (Map.Entry<Long, List<TsKvEntry>> entry : groupedTelemetry.entrySet()) { |
||||
|
ObjectNode element = json.createObjectNode(); |
||||
|
element.put("ts", entry.getKey()); |
||||
|
ObjectNode values = element.putObject("values"); |
||||
|
for (TsKvEntry tsKvEntry : entry.getValue()) { |
||||
|
addKvEntry(values, tsKvEntry); |
||||
|
} |
||||
|
result.add(element); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) throws Exception { |
||||
|
if (kvEntry.getDataType() == DataType.BOOLEAN) { |
||||
|
kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); |
||||
|
} else if (kvEntry.getDataType() == DataType.DOUBLE) { |
||||
|
kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); |
||||
|
} else if (kvEntry.getDataType() == DataType.LONG) { |
||||
|
kvEntry.getLongValue().ifPresent(value -> entityNode.put(kvEntry.getKey(), value)); |
||||
|
} else if (kvEntry.getDataType() == DataType.JSON) { |
||||
|
if (kvEntry.getJsonValue().isPresent()) { |
||||
|
entityNode.set(kvEntry.getKey(), json.readTree(kvEntry.getJsonValue().get())); |
||||
|
} |
||||
|
} else { |
||||
|
entityNode.put(kvEntry.getKey(), kvEntry.getValueAsString()); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,153 @@ |
|||||
|
/** |
||||
|
* 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.apiusage; |
||||
|
|
||||
|
import lombok.Getter; |
||||
|
import org.springframework.data.util.Pair; |
||||
|
import org.thingsboard.server.common.data.ApiFeature; |
||||
|
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
||||
|
import org.thingsboard.server.common.data.ApiUsageState; |
||||
|
import org.thingsboard.server.common.data.ApiUsageStateValue; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.tools.SchedulerUtils; |
||||
|
|
||||
|
import java.util.Arrays; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.HashSet; |
||||
|
import java.util.Map; |
||||
|
import java.util.Set; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
|
||||
|
public abstract class BaseApiUsageState { |
||||
|
private final Map<ApiUsageRecordKey, Long> currentCycleValues = new ConcurrentHashMap<>(); |
||||
|
private final Map<ApiUsageRecordKey, Long> currentHourValues = new ConcurrentHashMap<>(); |
||||
|
|
||||
|
@Getter |
||||
|
private final ApiUsageState apiUsageState; |
||||
|
@Getter |
||||
|
private volatile long currentCycleTs; |
||||
|
@Getter |
||||
|
private volatile long nextCycleTs; |
||||
|
@Getter |
||||
|
private volatile long currentHourTs; |
||||
|
|
||||
|
public BaseApiUsageState(ApiUsageState apiUsageState) { |
||||
|
this.apiUsageState = apiUsageState; |
||||
|
this.currentCycleTs = SchedulerUtils.getStartOfCurrentMonth(); |
||||
|
this.nextCycleTs = SchedulerUtils.getStartOfNextMonth(); |
||||
|
this.currentHourTs = SchedulerUtils.getStartOfCurrentHour(); |
||||
|
} |
||||
|
|
||||
|
public void put(ApiUsageRecordKey key, Long value) { |
||||
|
currentCycleValues.put(key, value); |
||||
|
} |
||||
|
|
||||
|
public void putHourly(ApiUsageRecordKey key, Long value) { |
||||
|
currentHourValues.put(key, value); |
||||
|
} |
||||
|
|
||||
|
public long add(ApiUsageRecordKey key, long value) { |
||||
|
long result = currentCycleValues.getOrDefault(key, 0L) + value; |
||||
|
currentCycleValues.put(key, result); |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
public long get(ApiUsageRecordKey key) { |
||||
|
return currentCycleValues.getOrDefault(key, 0L); |
||||
|
} |
||||
|
|
||||
|
public long addToHourly(ApiUsageRecordKey key, long value) { |
||||
|
long result = currentHourValues.getOrDefault(key, 0L) + value; |
||||
|
currentHourValues.put(key, result); |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
public void setHour(long currentHourTs) { |
||||
|
this.currentHourTs = currentHourTs; |
||||
|
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
||||
|
currentHourValues.put(key, 0L); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void setCycles(long currentCycleTs, long nextCycleTs) { |
||||
|
this.currentCycleTs = currentCycleTs; |
||||
|
this.nextCycleTs = nextCycleTs; |
||||
|
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
||||
|
currentCycleValues.put(key, 0L); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public ApiUsageStateValue getFeatureValue(ApiFeature feature) { |
||||
|
switch (feature) { |
||||
|
case TRANSPORT: |
||||
|
return apiUsageState.getTransportState(); |
||||
|
case RE: |
||||
|
return apiUsageState.getReExecState(); |
||||
|
case DB: |
||||
|
return apiUsageState.getDbStorageState(); |
||||
|
case JS: |
||||
|
return apiUsageState.getJsExecState(); |
||||
|
case EMAIL: |
||||
|
return apiUsageState.getEmailExecState(); |
||||
|
case SMS: |
||||
|
return apiUsageState.getSmsExecState(); |
||||
|
case ALARM: |
||||
|
return apiUsageState.getAlarmExecState(); |
||||
|
default: |
||||
|
return ApiUsageStateValue.ENABLED; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public boolean setFeatureValue(ApiFeature feature, ApiUsageStateValue value) { |
||||
|
ApiUsageStateValue currentValue = getFeatureValue(feature); |
||||
|
switch (feature) { |
||||
|
case TRANSPORT: |
||||
|
apiUsageState.setTransportState(value); |
||||
|
break; |
||||
|
case RE: |
||||
|
apiUsageState.setReExecState(value); |
||||
|
break; |
||||
|
case DB: |
||||
|
apiUsageState.setDbStorageState(value); |
||||
|
break; |
||||
|
case JS: |
||||
|
apiUsageState.setJsExecState(value); |
||||
|
break; |
||||
|
case EMAIL: |
||||
|
apiUsageState.setEmailExecState(value); |
||||
|
break; |
||||
|
case SMS: |
||||
|
apiUsageState.setSmsExecState(value); |
||||
|
break; |
||||
|
case ALARM: |
||||
|
apiUsageState.setAlarmExecState(value); |
||||
|
break; |
||||
|
} |
||||
|
return !currentValue.equals(value); |
||||
|
} |
||||
|
|
||||
|
public abstract EntityType getEntityType(); |
||||
|
|
||||
|
public TenantId getTenantId() { |
||||
|
return getApiUsageState().getTenantId(); |
||||
|
} |
||||
|
|
||||
|
public EntityId getEntityId() { |
||||
|
return getApiUsageState().getEntityId(); |
||||
|
} |
||||
|
} |
||||
@ -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.apiusage; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.ApiUsageState; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
|
||||
|
public class CustomerApiUsageState extends BaseApiUsageState { |
||||
|
public CustomerApiUsageState(ApiUsageState apiUsageState) { |
||||
|
super(apiUsageState); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public EntityType getEntityType() { |
||||
|
return EntityType.CUSTOMER; |
||||
|
} |
||||
|
} |
||||
@ -1,173 +0,0 @@ |
|||||
/** |
|
||||
* 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.lwm2m; |
|
||||
|
|
||||
|
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.eclipse.leshan.core.util.Hex; |
|
||||
import org.springframework.beans.factory.annotation.Autowired; |
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; |
|
||||
import org.springframework.stereotype.Service; |
|
||||
import org.thingsboard.server.common.data.lwm2m.ServerSecurityConfig; |
|
||||
import org.thingsboard.server.common.transport.lwm2m.LwM2MTransportConfigBootstrap; |
|
||||
import org.thingsboard.server.common.transport.lwm2m.LwM2MTransportConfigServer; |
|
||||
import org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode; |
|
||||
|
|
||||
import java.math.BigInteger; |
|
||||
import java.security.AlgorithmParameters; |
|
||||
import java.security.GeneralSecurityException; |
|
||||
import java.security.KeyFactory; |
|
||||
import java.security.KeyStoreException; |
|
||||
import java.security.PublicKey; |
|
||||
import java.security.cert.CertificateEncodingException; |
|
||||
import java.security.cert.X509Certificate; |
|
||||
import java.security.spec.ECGenParameterSpec; |
|
||||
import java.security.spec.ECParameterSpec; |
|
||||
import java.security.spec.ECPoint; |
|
||||
import java.security.spec.ECPublicKeySpec; |
|
||||
import java.security.spec.KeySpec; |
|
||||
|
|
||||
@Slf4j |
|
||||
@Service |
|
||||
@ConditionalOnExpression("('${service.type:null}'=='tb-transport' && '${transport.lwm2m.enabled:false}'=='true') || '${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core'") |
|
||||
public class LwM2MModelsRepository { |
|
||||
|
|
||||
private static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; |
|
||||
|
|
||||
@Autowired |
|
||||
LwM2MTransportConfigServer contextServer; |
|
||||
|
|
||||
|
|
||||
@Autowired |
|
||||
LwM2MTransportConfigBootstrap contextBootStrap; |
|
||||
|
|
||||
/** |
|
||||
* @param securityMode |
|
||||
* @param bootstrapServerIs |
|
||||
* @return ServerSecurityConfig more value is default: Important - port, host, publicKey |
|
||||
*/ |
|
||||
public ServerSecurityConfig getBootstrapSecurityInfo(String securityMode, boolean bootstrapServerIs) { |
|
||||
LwM2MSecurityMode lwM2MSecurityMode = LwM2MSecurityMode.fromSecurityMode(securityMode.toLowerCase()); |
|
||||
return getBootstrapServer(bootstrapServerIs, lwM2MSecurityMode); |
|
||||
} |
|
||||
|
|
||||
/** |
|
||||
* @param bootstrapServerIs |
|
||||
* @param mode |
|
||||
* @return ServerSecurityConfig more value is default: Important - port, host, publicKey |
|
||||
*/ |
|
||||
private ServerSecurityConfig getBootstrapServer(boolean bootstrapServerIs, LwM2MSecurityMode mode) { |
|
||||
ServerSecurityConfig bsServ = new ServerSecurityConfig(); |
|
||||
bsServ.setBootstrapServerIs(bootstrapServerIs); |
|
||||
if (bootstrapServerIs) { |
|
||||
bsServ.setServerId(contextBootStrap.getBootstrapServerId()); |
|
||||
switch (mode) { |
|
||||
case NO_SEC: |
|
||||
bsServ.setHost(contextBootStrap.getBootstrapHost()); |
|
||||
bsServ.setPort(contextBootStrap.getBootstrapPortNoSec()); |
|
||||
bsServ.setServerPublicKey(""); |
|
||||
break; |
|
||||
case PSK: |
|
||||
bsServ.setHost(contextBootStrap.getBootstrapHostSecurity()); |
|
||||
bsServ.setPort(contextBootStrap.getBootstrapPortSecurity()); |
|
||||
bsServ.setServerPublicKey(""); |
|
||||
break; |
|
||||
case RPK: |
|
||||
case X509: |
|
||||
bsServ.setHost(contextBootStrap.getBootstrapHostSecurity()); |
|
||||
bsServ.setPort(contextBootStrap.getBootstrapPortSecurity()); |
|
||||
bsServ.setServerPublicKey(getPublicKey (contextBootStrap.getBootstrapAlias(), this.contextBootStrap.getBootstrapPublicX(), this.contextBootStrap.getBootstrapPublicY())); |
|
||||
break; |
|
||||
default: |
|
||||
break; |
|
||||
} |
|
||||
} else { |
|
||||
bsServ.setServerId(contextServer.getServerId()); |
|
||||
switch (mode) { |
|
||||
case NO_SEC: |
|
||||
bsServ.setHost(contextServer.getServerHost()); |
|
||||
bsServ.setPort(contextServer.getServerPortNoSec()); |
|
||||
bsServ.setServerPublicKey(""); |
|
||||
break; |
|
||||
case PSK: |
|
||||
bsServ.setHost(contextServer.getServerHostSecurity()); |
|
||||
bsServ.setPort(contextServer.getServerPortSecurity()); |
|
||||
bsServ.setServerPublicKey(""); |
|
||||
break; |
|
||||
case RPK: |
|
||||
case X509: |
|
||||
bsServ.setHost(contextServer.getServerHostSecurity()); |
|
||||
bsServ.setPort(contextServer.getServerPortSecurity()); |
|
||||
bsServ.setServerPublicKey(getPublicKey (contextServer.getServerAlias(), this.contextServer.getServerPublicX(), this.contextServer.getServerPublicY())); |
|
||||
break; |
|
||||
default: |
|
||||
break; |
|
||||
} |
|
||||
} |
|
||||
return bsServ; |
|
||||
} |
|
||||
|
|
||||
private String getPublicKey (String alias, String publicServerX, String publicServerY) { |
|
||||
String publicKey = getServerPublicKeyX509(alias); |
|
||||
return publicKey != null ? publicKey : getRPKPublicKey(publicServerX, publicServerY); |
|
||||
} |
|
||||
|
|
||||
/** |
|
||||
* @param alias |
|
||||
* @return PublicKey format HexString or null |
|
||||
*/ |
|
||||
private String getServerPublicKeyX509(String alias) { |
|
||||
try { |
|
||||
X509Certificate serverCertificate = (X509Certificate) contextServer.getKeyStoreValue().getCertificate(alias); |
|
||||
return Hex.encodeHexString(serverCertificate.getEncoded()); |
|
||||
} catch (CertificateEncodingException | KeyStoreException e) { |
|
||||
e.printStackTrace(); |
|
||||
} |
|
||||
return null; |
|
||||
} |
|
||||
|
|
||||
/** |
|
||||
* @param publicServerX |
|
||||
* @param publicServerY |
|
||||
* @return PublicKey format HexString or null |
|
||||
*/ |
|
||||
private String getRPKPublicKey(String publicServerX, String publicServerY) { |
|
||||
try { |
|
||||
/** Get Elliptic Curve Parameter spec for secp256r1 */ |
|
||||
AlgorithmParameters algoParameters = AlgorithmParameters.getInstance("EC"); |
|
||||
algoParameters.init(new ECGenParameterSpec("secp256r1")); |
|
||||
ECParameterSpec parameterSpec = algoParameters.getParameterSpec(ECParameterSpec.class); |
|
||||
if (publicServerX != null && !publicServerX.isEmpty() && publicServerY != null && !publicServerY.isEmpty()) { |
|
||||
/** Get point values */ |
|
||||
byte[] publicX = Hex.decodeHex(publicServerX.toCharArray()); |
|
||||
byte[] publicY = Hex.decodeHex(publicServerY.toCharArray()); |
|
||||
/** Create key specs */ |
|
||||
KeySpec publicKeySpec = new ECPublicKeySpec(new ECPoint(new BigInteger(publicX), new BigInteger(publicY)), |
|
||||
parameterSpec); |
|
||||
/** Get keys */ |
|
||||
PublicKey publicKey = KeyFactory.getInstance("EC").generatePublic(publicKeySpec); |
|
||||
if (publicKey != null && publicKey.getEncoded().length > 0) { |
|
||||
return Hex.encodeHexString(publicKey.getEncoded()); |
|
||||
} |
|
||||
} |
|
||||
} catch (GeneralSecurityException | IllegalArgumentException e) { |
|
||||
log.error("[{}] Failed generate Server RPK for profile", e.getMessage()); |
|
||||
throw new RuntimeException(e); |
|
||||
} |
|
||||
return null; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@ -0,0 +1,110 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.lwm2m; |
||||
|
|
||||
|
|
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.eclipse.leshan.core.util.Hex; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.lwm2m.ServerSecurityConfig; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MSecureServerConfig; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportBootstrapConfig; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
||||
|
|
||||
|
import java.math.BigInteger; |
||||
|
import java.security.AlgorithmParameters; |
||||
|
import java.security.GeneralSecurityException; |
||||
|
import java.security.KeyFactory; |
||||
|
import java.security.KeyStoreException; |
||||
|
import java.security.PublicKey; |
||||
|
import java.security.cert.CertificateEncodingException; |
||||
|
import java.security.cert.X509Certificate; |
||||
|
import java.security.spec.ECGenParameterSpec; |
||||
|
import java.security.spec.ECParameterSpec; |
||||
|
import java.security.spec.ECPoint; |
||||
|
import java.security.spec.ECPublicKeySpec; |
||||
|
import java.security.spec.KeySpec; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Service |
||||
|
@RequiredArgsConstructor |
||||
|
@ConditionalOnExpression("('${service.type:null}'=='tb-transport' && '${transport.lwm2m.enabled:false}'=='true') || '${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core'") |
||||
|
public class LwM2MServerSecurityInfoRepository { |
||||
|
|
||||
|
private final LwM2MTransportServerConfig serverConfig; |
||||
|
private final LwM2MTransportBootstrapConfig bootstrapConfig; |
||||
|
|
||||
|
public ServerSecurityConfig getServerSecurityInfo(boolean bootstrapServer) { |
||||
|
ServerSecurityConfig result = getServerSecurityConfig(bootstrapServer ? bootstrapConfig : serverConfig); |
||||
|
result.setBootstrapServerIs(bootstrapServer); |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
private ServerSecurityConfig getServerSecurityConfig(LwM2MSecureServerConfig serverConfig) { |
||||
|
ServerSecurityConfig bsServ = new ServerSecurityConfig(); |
||||
|
bsServ.setServerId(serverConfig.getId()); |
||||
|
bsServ.setHost(serverConfig.getHost()); |
||||
|
bsServ.setPort(serverConfig.getPort()); |
||||
|
bsServ.setSecurityHost(serverConfig.getSecureHost()); |
||||
|
bsServ.setSecurityPort(serverConfig.getSecurePort()); |
||||
|
bsServ.setServerPublicKey(getPublicKey(serverConfig.getCertificateAlias(), this.serverConfig.getPublicX(), this.serverConfig.getPublicY())); |
||||
|
return bsServ; |
||||
|
} |
||||
|
|
||||
|
private String getPublicKey(String alias, String publicServerX, String publicServerY) { |
||||
|
String publicKey = getServerPublicKeyX509(alias); |
||||
|
return publicKey != null ? publicKey : getRPKPublicKey(publicServerX, publicServerY); |
||||
|
} |
||||
|
|
||||
|
private String getServerPublicKeyX509(String alias) { |
||||
|
try { |
||||
|
X509Certificate serverCertificate = (X509Certificate) serverConfig.getKeyStoreValue().getCertificate(alias); |
||||
|
return Hex.encodeHexString(serverCertificate.getEncoded()); |
||||
|
} catch (CertificateEncodingException | KeyStoreException e) { |
||||
|
e.printStackTrace(); |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
private String getRPKPublicKey(String publicServerX, String publicServerY) { |
||||
|
try { |
||||
|
/** Get Elliptic Curve Parameter spec for secp256r1 */ |
||||
|
AlgorithmParameters algoParameters = AlgorithmParameters.getInstance("EC"); |
||||
|
algoParameters.init(new ECGenParameterSpec("secp256r1")); |
||||
|
ECParameterSpec parameterSpec = algoParameters.getParameterSpec(ECParameterSpec.class); |
||||
|
if (publicServerX != null && !publicServerX.isEmpty() && publicServerY != null && !publicServerY.isEmpty()) { |
||||
|
/** Get point values */ |
||||
|
byte[] publicX = Hex.decodeHex(publicServerX.toCharArray()); |
||||
|
byte[] publicY = Hex.decodeHex(publicServerY.toCharArray()); |
||||
|
/** Create key specs */ |
||||
|
KeySpec publicKeySpec = new ECPublicKeySpec(new ECPoint(new BigInteger(publicX), new BigInteger(publicY)), |
||||
|
parameterSpec); |
||||
|
/** Get keys */ |
||||
|
PublicKey publicKey = KeyFactory.getInstance("EC").generatePublic(publicKeySpec); |
||||
|
if (publicKey != null && publicKey.getEncoded().length > 0) { |
||||
|
return Hex.encodeHexString(publicKey.getEncoded()); |
||||
|
} |
||||
|
} |
||||
|
} catch (GeneralSecurityException | IllegalArgumentException e) { |
||||
|
log.error("[{}] Failed generate Server RPK for profile", e.getMessage()); |
||||
|
throw new RuntimeException(e); |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
@ -0,0 +1,352 @@ |
|||||
|
/** |
||||
|
* 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.ota; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FutureCallback; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; |
||||
|
import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; |
||||
|
import org.thingsboard.server.common.data.DataConstants; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.DeviceProfile; |
||||
|
import org.thingsboard.server.common.data.OtaPackageInfo; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.OtaPackageId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKey; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.LongDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
||||
|
import org.thingsboard.server.common.data.ota.OtaPackageType; |
||||
|
import org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus; |
||||
|
import org.thingsboard.server.common.data.ota.OtaPackageUtil; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
||||
|
import org.thingsboard.server.dao.device.DeviceProfileService; |
||||
|
import org.thingsboard.server.dao.device.DeviceService; |
||||
|
import org.thingsboard.server.dao.ota.OtaPackageService; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg; |
||||
|
import org.thingsboard.server.queue.TbQueueProducer; |
||||
|
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
||||
|
import org.thingsboard.server.queue.provider.TbCoreQueueFactory; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.queue.TbClusterService; |
||||
|
|
||||
|
import javax.annotation.Nullable; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.Collections; |
||||
|
import java.util.HashSet; |
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
import java.util.UUID; |
||||
|
import java.util.function.Consumer; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.CHECKSUM; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.CHECKSUM_ALGORITHM; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.SIZE; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.STATE; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.TITLE; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.TS; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.URL; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageKey.VERSION; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageType.SOFTWARE; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageUtil.getAttributeKey; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageUtil.getTargetTelemetryKey; |
||||
|
import static org.thingsboard.server.common.data.ota.OtaPackageUtil.getTelemetryKey; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Service |
||||
|
@TbCoreComponent |
||||
|
public class DefaultOtaPackageStateService implements OtaPackageStateService { |
||||
|
|
||||
|
private final TbClusterService tbClusterService; |
||||
|
private final OtaPackageService otaPackageService; |
||||
|
private final DeviceService deviceService; |
||||
|
private final DeviceProfileService deviceProfileService; |
||||
|
private final RuleEngineTelemetryService telemetryService; |
||||
|
private final TbQueueProducer<TbProtoQueueMsg<ToOtaPackageStateServiceMsg>> otaPackageStateMsgProducer; |
||||
|
|
||||
|
public DefaultOtaPackageStateService(TbClusterService tbClusterService, OtaPackageService otaPackageService, |
||||
|
DeviceService deviceService, |
||||
|
DeviceProfileService deviceProfileService, |
||||
|
RuleEngineTelemetryService telemetryService, |
||||
|
TbCoreQueueFactory coreQueueFactory) { |
||||
|
this.tbClusterService = tbClusterService; |
||||
|
this.otaPackageService = otaPackageService; |
||||
|
this.deviceService = deviceService; |
||||
|
this.deviceProfileService = deviceProfileService; |
||||
|
this.telemetryService = telemetryService; |
||||
|
this.otaPackageStateMsgProducer = coreQueueFactory.createToOtaPackageStateServiceMsgProducer(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void update(Device device, Device oldDevice) { |
||||
|
updateFirmware(device, oldDevice); |
||||
|
updateSoftware(device, oldDevice); |
||||
|
} |
||||
|
|
||||
|
private void updateFirmware(Device device, Device oldDevice) { |
||||
|
OtaPackageId newFirmwareId = device.getFirmwareId(); |
||||
|
if (newFirmwareId == null) { |
||||
|
DeviceProfile newDeviceProfile = deviceProfileService.findDeviceProfileById(device.getTenantId(), device.getDeviceProfileId()); |
||||
|
newFirmwareId = newDeviceProfile.getFirmwareId(); |
||||
|
} |
||||
|
if (oldDevice != null) { |
||||
|
if (newFirmwareId != null) { |
||||
|
OtaPackageId oldFirmwareId = oldDevice.getFirmwareId(); |
||||
|
if (oldFirmwareId == null) { |
||||
|
DeviceProfile oldDeviceProfile = deviceProfileService.findDeviceProfileById(oldDevice.getTenantId(), oldDevice.getDeviceProfileId()); |
||||
|
oldFirmwareId = oldDeviceProfile.getFirmwareId(); |
||||
|
} |
||||
|
if (!newFirmwareId.equals(oldFirmwareId)) { |
||||
|
// Device was updated and new firmware is different from previous firmware.
|
||||
|
send(device.getTenantId(), device.getId(), newFirmwareId, System.currentTimeMillis(), FIRMWARE); |
||||
|
} |
||||
|
} else { |
||||
|
// Device was updated and new firmware is not set.
|
||||
|
remove(device, FIRMWARE); |
||||
|
} |
||||
|
} else if (newFirmwareId != null) { |
||||
|
// Device was created and firmware is defined.
|
||||
|
send(device.getTenantId(), device.getId(), newFirmwareId, System.currentTimeMillis(), FIRMWARE); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void updateSoftware(Device device, Device oldDevice) { |
||||
|
OtaPackageId newSoftwareId = device.getSoftwareId(); |
||||
|
if (newSoftwareId == null) { |
||||
|
DeviceProfile newDeviceProfile = deviceProfileService.findDeviceProfileById(device.getTenantId(), device.getDeviceProfileId()); |
||||
|
newSoftwareId = newDeviceProfile.getSoftwareId(); |
||||
|
} |
||||
|
if (oldDevice != null) { |
||||
|
if (newSoftwareId != null) { |
||||
|
OtaPackageId oldSoftwareId = oldDevice.getSoftwareId(); |
||||
|
if (oldSoftwareId == null) { |
||||
|
DeviceProfile oldDeviceProfile = deviceProfileService.findDeviceProfileById(oldDevice.getTenantId(), oldDevice.getDeviceProfileId()); |
||||
|
oldSoftwareId = oldDeviceProfile.getSoftwareId(); |
||||
|
} |
||||
|
if (!newSoftwareId.equals(oldSoftwareId)) { |
||||
|
// Device was updated and new firmware is different from previous firmware.
|
||||
|
send(device.getTenantId(), device.getId(), newSoftwareId, System.currentTimeMillis(), SOFTWARE); |
||||
|
} |
||||
|
} else { |
||||
|
// Device was updated and new firmware is not set.
|
||||
|
remove(device, SOFTWARE); |
||||
|
} |
||||
|
} else if (newSoftwareId != null) { |
||||
|
// Device was created and firmware is defined.
|
||||
|
send(device.getTenantId(), device.getId(), newSoftwareId, System.currentTimeMillis(), SOFTWARE); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void update(DeviceProfile deviceProfile, boolean isFirmwareChanged, boolean isSoftwareChanged) { |
||||
|
TenantId tenantId = deviceProfile.getTenantId(); |
||||
|
|
||||
|
if (isFirmwareChanged) { |
||||
|
update(tenantId, deviceProfile, FIRMWARE); |
||||
|
} |
||||
|
if (isSoftwareChanged) { |
||||
|
update(tenantId, deviceProfile, SOFTWARE); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void update(TenantId tenantId, DeviceProfile deviceProfile, OtaPackageType otaPackageType) { |
||||
|
Consumer<Device> updateConsumer; |
||||
|
|
||||
|
if (deviceProfile.getFirmwareId() != null) { |
||||
|
long ts = System.currentTimeMillis(); |
||||
|
updateConsumer = d -> send(d.getTenantId(), d.getId(), deviceProfile.getFirmwareId(), ts, otaPackageType); |
||||
|
} else { |
||||
|
updateConsumer = d -> remove(d, otaPackageType); |
||||
|
} |
||||
|
|
||||
|
PageLink pageLink = new PageLink(100); |
||||
|
PageData<Device> pageData; |
||||
|
do { |
||||
|
pageData = deviceService.findDevicesByTenantIdAndTypeAndEmptyOtaPackage(tenantId, deviceProfile.getId(), otaPackageType, pageLink); |
||||
|
pageData.getData().forEach(updateConsumer); |
||||
|
|
||||
|
if (pageData.hasNext()) { |
||||
|
pageLink = pageLink.nextPageLink(); |
||||
|
} |
||||
|
} while (pageData.hasNext()); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean process(ToOtaPackageStateServiceMsg msg) { |
||||
|
boolean isSuccess = false; |
||||
|
OtaPackageId targetOtaPackageId = new OtaPackageId(new UUID(msg.getOtaPackageIdMSB(), msg.getOtaPackageIdLSB())); |
||||
|
DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB())); |
||||
|
TenantId tenantId = new TenantId(new UUID(msg.getTenantIdMSB(), msg.getTenantIdLSB())); |
||||
|
OtaPackageType firmwareType = OtaPackageType.valueOf(msg.getType()); |
||||
|
long ts = msg.getTs(); |
||||
|
|
||||
|
Device device = deviceService.findDeviceById(tenantId, deviceId); |
||||
|
if (device == null) { |
||||
|
log.warn("[{}] [{}] Device was removed during firmware update msg was queued!", tenantId, deviceId); |
||||
|
} else { |
||||
|
OtaPackageId currentOtaPackageId = OtaPackageUtil.getOtaPackageId(device, firmwareType); |
||||
|
if (currentOtaPackageId == null) { |
||||
|
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(tenantId, device.getDeviceProfileId()); |
||||
|
currentOtaPackageId = OtaPackageUtil.getOtaPackageId(deviceProfile, firmwareType); |
||||
|
} |
||||
|
|
||||
|
if (targetOtaPackageId.equals(currentOtaPackageId)) { |
||||
|
update(device, otaPackageService.findOtaPackageInfoById(device.getTenantId(), targetOtaPackageId), ts); |
||||
|
isSuccess = true; |
||||
|
} else { |
||||
|
log.warn("[{}] [{}] Can`t update firmware for the device, target firmwareId: [{}], current firmwareId: [{}]!", tenantId, deviceId, targetOtaPackageId, currentOtaPackageId); |
||||
|
} |
||||
|
} |
||||
|
return isSuccess; |
||||
|
} |
||||
|
|
||||
|
private void send(TenantId tenantId, DeviceId deviceId, OtaPackageId firmwareId, long ts, OtaPackageType firmwareType) { |
||||
|
ToOtaPackageStateServiceMsg msg = ToOtaPackageStateServiceMsg.newBuilder() |
||||
|
.setTenantIdMSB(tenantId.getId().getMostSignificantBits()) |
||||
|
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) |
||||
|
.setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) |
||||
|
.setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) |
||||
|
.setOtaPackageIdMSB(firmwareId.getId().getMostSignificantBits()) |
||||
|
.setOtaPackageIdLSB(firmwareId.getId().getLeastSignificantBits()) |
||||
|
.setType(firmwareType.name()) |
||||
|
.setTs(ts) |
||||
|
.build(); |
||||
|
|
||||
|
OtaPackageInfo firmware = otaPackageService.findOtaPackageInfoById(tenantId, firmwareId); |
||||
|
if (firmware == null) { |
||||
|
log.warn("[{}] Failed to send firmware update because firmware was already deleted", firmwareId); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
TopicPartitionInfo tpi = new TopicPartitionInfo(otaPackageStateMsgProducer.getDefaultTopic(), null, null, false); |
||||
|
otaPackageStateMsgProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); |
||||
|
|
||||
|
List<TsKvEntry> telemetry = new ArrayList<>(); |
||||
|
telemetry.add(new BasicTsKvEntry(ts, new StringDataEntry(getTargetTelemetryKey(firmware.getType(), TITLE), firmware.getTitle()))); |
||||
|
telemetry.add(new BasicTsKvEntry(ts, new StringDataEntry(getTargetTelemetryKey(firmware.getType(), VERSION), firmware.getVersion()))); |
||||
|
telemetry.add(new BasicTsKvEntry(ts, new LongDataEntry(getTargetTelemetryKey(firmware.getType(), TS), ts))); |
||||
|
telemetry.add(new BasicTsKvEntry(ts, new StringDataEntry(getTelemetryKey(firmware.getType(), STATE), OtaPackageUpdateStatus.QUEUED.name()))); |
||||
|
|
||||
|
telemetryService.saveAndNotify(tenantId, deviceId, telemetry, new FutureCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(@Nullable Void tmp) { |
||||
|
log.trace("[{}] Success save firmware status!", deviceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Throwable t) { |
||||
|
log.error("[{}] Failed to save firmware status!", deviceId, t); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
|
||||
|
private void update(Device device, OtaPackageInfo otaPackage, long ts) { |
||||
|
TenantId tenantId = device.getTenantId(); |
||||
|
DeviceId deviceId = device.getId(); |
||||
|
OtaPackageType otaPackageType = otaPackage.getType(); |
||||
|
|
||||
|
BasicTsKvEntry status = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(getTelemetryKey(otaPackageType, STATE), OtaPackageUpdateStatus.INITIATED.name())); |
||||
|
|
||||
|
telemetryService.saveAndNotify(tenantId, deviceId, Collections.singletonList(status), new FutureCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(@Nullable Void tmp) { |
||||
|
log.trace("[{}] Success save telemetry with target firmware for device!", deviceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Throwable t) { |
||||
|
log.error("[{}] Failed to save telemetry with target firmware for device!", deviceId, t); |
||||
|
} |
||||
|
}); |
||||
|
|
||||
|
List<AttributeKvEntry> attributes = new ArrayList<>(); |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, TITLE), otaPackage.getTitle()))); |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, VERSION), otaPackage.getVersion()))); |
||||
|
if (otaPackage.hasUrl()) { |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, URL), otaPackage.getUrl()))); |
||||
|
List<String> attrToRemove = new ArrayList<>(); |
||||
|
|
||||
|
if (otaPackage.getDataSize() == null) { |
||||
|
attrToRemove.add(getAttributeKey(otaPackageType, SIZE)); |
||||
|
} else { |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new LongDataEntry(getAttributeKey(otaPackageType, SIZE), otaPackage.getDataSize()))); |
||||
|
} |
||||
|
|
||||
|
if (otaPackage.getChecksumAlgorithm() != null) { |
||||
|
attrToRemove.add(getAttributeKey(otaPackageType, CHECKSUM_ALGORITHM)); |
||||
|
} else { |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM_ALGORITHM), otaPackage.getChecksumAlgorithm().name()))); |
||||
|
} |
||||
|
|
||||
|
if (StringUtils.isEmpty(otaPackage.getChecksum())) { |
||||
|
attrToRemove.add(getAttributeKey(otaPackageType, CHECKSUM)); |
||||
|
} else { |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM), otaPackage.getChecksum()))); |
||||
|
} |
||||
|
|
||||
|
remove(device, otaPackageType, attrToRemove); |
||||
|
} else { |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new LongDataEntry(getAttributeKey(otaPackageType, SIZE), otaPackage.getDataSize()))); |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM_ALGORITHM), otaPackage.getChecksumAlgorithm().name()))); |
||||
|
attributes.add(new BaseAttributeKvEntry(ts, new StringDataEntry(getAttributeKey(otaPackageType, CHECKSUM), otaPackage.getChecksum()))); |
||||
|
remove(device, otaPackageType, Collections.singletonList(getAttributeKey(otaPackageType, URL))); |
||||
|
} |
||||
|
|
||||
|
telemetryService.saveAndNotify(tenantId, deviceId, DataConstants.SHARED_SCOPE, attributes, new FutureCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(@Nullable Void tmp) { |
||||
|
log.trace("[{}] Success save attributes with target firmware!", deviceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Throwable t) { |
||||
|
log.error("[{}] Failed to save attributes with target firmware!", deviceId, t); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
private void remove(Device device, OtaPackageType otaPackageType) { |
||||
|
remove(device, otaPackageType, OtaPackageUtil.getAttributeKeys(otaPackageType)); |
||||
|
} |
||||
|
|
||||
|
private void remove(Device device, OtaPackageType otaPackageType, List<String> attributesKeys) { |
||||
|
telemetryService.deleteAndNotify(device.getTenantId(), device.getId(), DataConstants.SHARED_SCOPE, attributesKeys, |
||||
|
new FutureCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(@Nullable Void tmp) { |
||||
|
log.trace("[{}] Success remove target {} attributes!", device.getId(), otaPackageType); |
||||
|
Set<AttributeKey> keysToNotify = new HashSet<>(); |
||||
|
attributesKeys.forEach(key -> keysToNotify.add(new AttributeKey(DataConstants.SHARED_SCOPE, key))); |
||||
|
tbClusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete(device.getTenantId(), device.getId(), keysToNotify), null); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Throwable t) { |
||||
|
log.error("[{}] Failed to remove target {} attributes!", device.getId(), otaPackageType, t); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -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.ota; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.DeviceProfile; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg; |
||||
|
|
||||
|
public interface OtaPackageStateService { |
||||
|
|
||||
|
void update(Device device, Device oldDevice); |
||||
|
|
||||
|
void update(DeviceProfile deviceProfile, boolean isFirmwareChanged, boolean isSoftwareChanged); |
||||
|
|
||||
|
boolean process(ToOtaPackageStateServiceMsg msg); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,218 @@ |
|||||
|
/** |
||||
|
* 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.resource; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.apache.commons.lang3.StringUtils; |
||||
|
import org.eclipse.leshan.core.model.DDFFileParser; |
||||
|
import org.eclipse.leshan.core.model.DefaultDDFFileValidator; |
||||
|
import org.eclipse.leshan.core.model.InvalidDDFFileException; |
||||
|
import org.eclipse.leshan.core.model.ObjectModel; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.ResourceType; |
||||
|
import org.thingsboard.server.common.data.TbResource; |
||||
|
import org.thingsboard.server.common.data.TbResourceInfo; |
||||
|
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; |
||||
|
import org.thingsboard.server.common.data.exception.ThingsboardException; |
||||
|
import org.thingsboard.server.common.data.id.TbResourceId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.lwm2m.LwM2mInstance; |
||||
|
import org.thingsboard.server.common.data.lwm2m.LwM2mObject; |
||||
|
import org.thingsboard.server.common.data.lwm2m.LwM2mResourceObserve; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.dao.exception.DataValidationException; |
||||
|
import org.thingsboard.server.dao.resource.ResourceService; |
||||
|
|
||||
|
import java.io.ByteArrayInputStream; |
||||
|
import java.io.IOException; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.Base64; |
||||
|
import java.util.Comparator; |
||||
|
import java.util.List; |
||||
|
import java.util.stream.Collectors; |
||||
|
import java.util.stream.Stream; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_KEY; |
||||
|
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_SEARCH_TEXT; |
||||
|
import static org.thingsboard.server.dao.device.DeviceServiceImpl.INCORRECT_TENANT_ID; |
||||
|
import static org.thingsboard.server.dao.service.Validator.validateId; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Service |
||||
|
public class DefaultTbResourceService implements TbResourceService { |
||||
|
|
||||
|
private final ResourceService resourceService; |
||||
|
private final DDFFileParser ddfFileParser; |
||||
|
|
||||
|
public DefaultTbResourceService(ResourceService resourceService) { |
||||
|
this.resourceService = resourceService; |
||||
|
this.ddfFileParser = new DDFFileParser(new DefaultDDFFileValidator()); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public TbResource saveResource(TbResource resource) throws ThingsboardException { |
||||
|
log.trace("Executing saveResource [{}]", resource); |
||||
|
if (StringUtils.isEmpty(resource.getData())) { |
||||
|
throw new DataValidationException("Resource data should be specified!"); |
||||
|
} |
||||
|
if (ResourceType.LWM2M_MODEL.equals(resource.getResourceType())) { |
||||
|
try { |
||||
|
List<ObjectModel> objectModels = |
||||
|
ddfFileParser.parse(new ByteArrayInputStream(Base64.getDecoder().decode(resource.getData())), resource.getSearchText()); |
||||
|
if (!objectModels.isEmpty()) { |
||||
|
ObjectModel objectModel = objectModels.get(0); |
||||
|
|
||||
|
String resourceKey = objectModel.id + LWM2M_SEPARATOR_KEY + objectModel.version; |
||||
|
String name = objectModel.name; |
||||
|
resource.setResourceKey(resourceKey); |
||||
|
if (resource.getId() == null) { |
||||
|
resource.setTitle(name + " id=" + objectModel.id + " v" + objectModel.version); |
||||
|
} |
||||
|
resource.setSearchText(resourceKey + LWM2M_SEPARATOR_SEARCH_TEXT + name); |
||||
|
} else { |
||||
|
throw new DataValidationException(String.format("Could not parse the XML of objectModel with name %s", resource.getSearchText())); |
||||
|
} |
||||
|
} catch (InvalidDDFFileException e) { |
||||
|
log.error("Failed to parse file {}", resource.getFileName(), e); |
||||
|
throw new DataValidationException("Failed to parse file " + resource.getFileName()); |
||||
|
} catch (IOException e) { |
||||
|
throw new ThingsboardException(e, ThingsboardErrorCode.GENERAL); |
||||
|
} |
||||
|
if (resource.getResourceType().equals(ResourceType.LWM2M_MODEL) && toLwM2mObject(resource, true) == null) { |
||||
|
throw new DataValidationException(String.format("Could not parse the XML of objectModel with name %s", resource.getSearchText())); |
||||
|
} |
||||
|
} else { |
||||
|
resource.setResourceKey(resource.getFileName()); |
||||
|
} |
||||
|
|
||||
|
return resourceService.saveResource(resource); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public TbResource getResource(TenantId tenantId, ResourceType resourceType, String resourceId) { |
||||
|
return resourceService.getResource(tenantId, resourceType, resourceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public TbResource findResourceById(TenantId tenantId, TbResourceId resourceId) { |
||||
|
return resourceService.findResourceById(tenantId, resourceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public TbResourceInfo findResourceInfoById(TenantId tenantId, TbResourceId resourceId) { |
||||
|
return resourceService.findResourceInfoById(tenantId, resourceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public PageData<TbResourceInfo> findAllTenantResourcesByTenantId(TenantId tenantId, PageLink pageLink) { |
||||
|
return resourceService.findAllTenantResourcesByTenantId(tenantId, pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public PageData<TbResourceInfo> findTenantResourcesByTenantId(TenantId tenantId, PageLink pageLink) { |
||||
|
return resourceService.findTenantResourcesByTenantId(tenantId, pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public List<LwM2mObject> findLwM2mObject(TenantId tenantId, String sortOrder, String sortProperty, String[] objectIds) { |
||||
|
log.trace("Executing findByTenantId [{}]", tenantId); |
||||
|
validateId(tenantId, INCORRECT_TENANT_ID + tenantId); |
||||
|
List<TbResource> resources = resourceService.findTenantResourcesByResourceTypeAndObjectIds(tenantId, ResourceType.LWM2M_MODEL, |
||||
|
objectIds); |
||||
|
return resources.stream() |
||||
|
.flatMap(s -> Stream.ofNullable(toLwM2mObject(s, false))) |
||||
|
.sorted(getComparator(sortProperty, sortOrder)) |
||||
|
.collect(Collectors.toList()); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public List<LwM2mObject> findLwM2mObjectPage(TenantId tenantId, String sortProperty, String sortOrder, PageLink pageLink) { |
||||
|
log.trace("Executing findByTenantId [{}]", tenantId); |
||||
|
validateId(tenantId, INCORRECT_TENANT_ID + tenantId); |
||||
|
PageData<TbResource> resourcePageData = resourceService.findTenantResourcesByResourceTypeAndPageLink(tenantId, ResourceType.LWM2M_MODEL, pageLink); |
||||
|
return resourcePageData.getData().stream() |
||||
|
.flatMap(s -> Stream.ofNullable(toLwM2mObject(s, false))) |
||||
|
.sorted(getComparator(sortProperty, sortOrder)) |
||||
|
.collect(Collectors.toList()); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void deleteResource(TenantId tenantId, TbResourceId resourceId) { |
||||
|
resourceService.deleteResource(tenantId, resourceId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void deleteResourcesByTenantId(TenantId tenantId) { |
||||
|
resourceService.deleteResourcesByTenantId(tenantId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public long sumDataSizeByTenantId(TenantId tenantId) { |
||||
|
return resourceService.sumDataSizeByTenantId(tenantId); |
||||
|
} |
||||
|
|
||||
|
private Comparator<? super LwM2mObject> getComparator(String sortProperty, String sortOrder) { |
||||
|
Comparator<LwM2mObject> comparator; |
||||
|
if ("name".equals(sortProperty)) { |
||||
|
comparator = Comparator.comparing(LwM2mObject::getName); |
||||
|
} else { |
||||
|
comparator = Comparator.comparingLong(LwM2mObject::getId); |
||||
|
} |
||||
|
return "DESC".equals(sortOrder) ? comparator.reversed() : comparator; |
||||
|
} |
||||
|
|
||||
|
private LwM2mObject toLwM2mObject(TbResource resource, boolean isSave) { |
||||
|
try { |
||||
|
DDFFileParser ddfFileParser = new DDFFileParser(new DefaultDDFFileValidator()); |
||||
|
List<ObjectModel> objectModels = |
||||
|
ddfFileParser.parse(new ByteArrayInputStream(Base64.getDecoder().decode(resource.getData())), resource.getSearchText()); |
||||
|
if (objectModels.size() == 0) { |
||||
|
return null; |
||||
|
} else { |
||||
|
ObjectModel obj = objectModels.get(0); |
||||
|
LwM2mObject lwM2mObject = new LwM2mObject(); |
||||
|
lwM2mObject.setId(obj.id); |
||||
|
lwM2mObject.setKeyId(resource.getResourceKey()); |
||||
|
lwM2mObject.setName(obj.name); |
||||
|
lwM2mObject.setMultiple(obj.multiple); |
||||
|
lwM2mObject.setMandatory(obj.mandatory); |
||||
|
LwM2mInstance instance = new LwM2mInstance(); |
||||
|
instance.setId(0); |
||||
|
List<LwM2mResourceObserve> resources = new ArrayList<>(); |
||||
|
obj.resources.forEach((k, v) -> { |
||||
|
if (isSave) { |
||||
|
LwM2mResourceObserve lwM2MResourceObserve = new LwM2mResourceObserve(k, v.name, false, false, false); |
||||
|
resources.add(lwM2MResourceObserve); |
||||
|
} else if (v.operations.isReadable()) { |
||||
|
LwM2mResourceObserve lwM2MResourceObserve = new LwM2mResourceObserve(k, v.name, false, false, false); |
||||
|
resources.add(lwM2MResourceObserve); |
||||
|
} |
||||
|
}); |
||||
|
if (isSave || resources.size() > 0) { |
||||
|
instance.setResources(resources.toArray(LwM2mResourceObserve[]::new)); |
||||
|
lwM2mObject.setInstances(new LwM2mInstance[]{instance}); |
||||
|
return lwM2mObject; |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
} catch (IOException | InvalidDDFFileException e) { |
||||
|
log.error("Could not parse the XML of objectModel with name [{}]", resource.getSearchText(), e); |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -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.rpc; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.RpcId; |
||||
|
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.rpc.Rpc; |
||||
|
import org.thingsboard.server.common.data.rpc.RpcStatus; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
import org.thingsboard.server.dao.rpc.RpcService; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.queue.TbClusterService; |
||||
|
|
||||
|
@TbCoreComponent |
||||
|
@Service |
||||
|
@RequiredArgsConstructor |
||||
|
@Slf4j |
||||
|
public class TbRpcService { |
||||
|
private final RpcService rpcService; |
||||
|
private final TbClusterService tbClusterService; |
||||
|
|
||||
|
public Rpc save(TenantId tenantId, Rpc rpc) { |
||||
|
Rpc saved = rpcService.save(rpc); |
||||
|
pushRpcMsgToRuleEngine(tenantId, saved); |
||||
|
return saved; |
||||
|
} |
||||
|
|
||||
|
public void save(TenantId tenantId, RpcId rpcId, RpcStatus newStatus, JsonNode response) { |
||||
|
Rpc foundRpc = rpcService.findById(tenantId, rpcId); |
||||
|
if (foundRpc != null) { |
||||
|
foundRpc.setStatus(newStatus); |
||||
|
if (response != null) { |
||||
|
foundRpc.setResponse(response); |
||||
|
} |
||||
|
Rpc saved = rpcService.save(foundRpc); |
||||
|
pushRpcMsgToRuleEngine(tenantId, saved); |
||||
|
} else { |
||||
|
log.warn("[{}] Failed to update RPC status because RPC was already deleted", rpcId); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) { |
||||
|
TbMsg msg = TbMsg.newMsg("RPC_" + rpc.getStatus().name(), rpc.getDeviceId(), TbMsgMetaData.EMPTY, JacksonUtil.toString(rpc)); |
||||
|
tbClusterService.pushMsgToRuleEngine(tenantId, rpc.getId(), msg, null); |
||||
|
} |
||||
|
|
||||
|
public Rpc findRpcById(TenantId tenantId, RpcId rpcId) { |
||||
|
return rpcService.findById(tenantId, rpcId); |
||||
|
} |
||||
|
|
||||
|
public PageData<Rpc> findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId deviceId, RpcStatus rpcStatus, PageLink pageLink) { |
||||
|
return rpcService.findAllByDeviceIdAndStatus(tenantId, deviceId, rpcStatus, pageLink); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,101 @@ |
|||||
|
/** |
||||
|
* 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.security.auth.oauth2; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.security.oauth2.client.authentication.OAuth2AuthenticationToken; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.springframework.util.LinkedMultiValueMap; |
||||
|
import org.springframework.util.MultiValueMap; |
||||
|
import org.springframework.util.StringUtils; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.oauth2.OAuth2MapperConfig; |
||||
|
import org.thingsboard.server.common.data.oauth2.OAuth2Registration; |
||||
|
import org.thingsboard.server.dao.oauth2.OAuth2User; |
||||
|
import org.thingsboard.server.service.security.model.SecurityUser; |
||||
|
|
||||
|
import javax.servlet.http.HttpServletRequest; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.Map; |
||||
|
|
||||
|
@Service(value = "appleOAuth2ClientMapper") |
||||
|
@Slf4j |
||||
|
public class AppleOAuth2ClientMapper extends AbstractOAuth2ClientMapper implements OAuth2ClientMapper { |
||||
|
|
||||
|
private static final String USER = "user"; |
||||
|
private static final String NAME = "name"; |
||||
|
private static final String FIRST_NAME = "firstName"; |
||||
|
private static final String LAST_NAME = "lastName"; |
||||
|
private static final String EMAIL = "email"; |
||||
|
|
||||
|
@Override |
||||
|
public SecurityUser getOrCreateUserByClientPrincipal(HttpServletRequest request, OAuth2AuthenticationToken token, String providerAccessToken, OAuth2Registration registration) { |
||||
|
OAuth2MapperConfig config = registration.getMapperConfig(); |
||||
|
Map<String, Object> attributes = updateAttributesFromRequestParams(request, token.getPrincipal().getAttributes()); |
||||
|
String email = BasicMapperUtils.getStringAttributeByKey(attributes, config.getBasic().getEmailAttributeKey()); |
||||
|
OAuth2User oauth2User = BasicMapperUtils.getOAuth2User(email, attributes, config); |
||||
|
|
||||
|
return getOrCreateSecurityUserFromOAuth2User(oauth2User, registration); |
||||
|
} |
||||
|
|
||||
|
private static Map<String, Object> updateAttributesFromRequestParams(HttpServletRequest request, Map<String, Object> attributes) { |
||||
|
Map<String, Object> updated = attributes; |
||||
|
MultiValueMap<String, String> params = toMultiMap(request.getParameterMap()); |
||||
|
String userValue = params.getFirst(USER); |
||||
|
if (StringUtils.hasText(userValue)) { |
||||
|
JsonNode user = null; |
||||
|
try { |
||||
|
user = JacksonUtil.toJsonNode(userValue); |
||||
|
} catch (Exception e) {} |
||||
|
if (user != null) { |
||||
|
updated = new HashMap<>(attributes); |
||||
|
if (user.has(NAME)) { |
||||
|
JsonNode name = user.get(NAME); |
||||
|
if (name.isObject()) { |
||||
|
JsonNode firstName = name.get(FIRST_NAME); |
||||
|
if (firstName != null && firstName.isTextual()) { |
||||
|
updated.put(FIRST_NAME, firstName.asText()); |
||||
|
} |
||||
|
JsonNode lastName = name.get(LAST_NAME); |
||||
|
if (lastName != null && lastName.isTextual()) { |
||||
|
updated.put(LAST_NAME, lastName.asText()); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
if (user.has(EMAIL)) { |
||||
|
JsonNode email = user.get(EMAIL); |
||||
|
if (email != null && email.isTextual()) { |
||||
|
updated.put(EMAIL, email.asText()); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return updated; |
||||
|
} |
||||
|
|
||||
|
private static MultiValueMap<String, String> toMultiMap(Map<String, String[]> map) { |
||||
|
MultiValueMap<String, String> params = new LinkedMultiValueMap<>(map.size()); |
||||
|
map.forEach((key, values) -> { |
||||
|
if (values.length > 0) { |
||||
|
for (String value : values) { |
||||
|
params.add(key, value); |
||||
|
} |
||||
|
} |
||||
|
}); |
||||
|
return params; |
||||
|
} |
||||
|
} |
||||
@ -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.service.security.auth.oauth2; |
||||
|
|
||||
|
public interface TbOAuth2ParameterNames { |
||||
|
|
||||
|
String CALLBACK_URL_SCHEME = "callback_url_scheme"; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,69 @@ |
|||||
|
/** |
||||
|
* 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.security.model.token; |
||||
|
|
||||
|
import io.jsonwebtoken.Claims; |
||||
|
import io.jsonwebtoken.ExpiredJwtException; |
||||
|
import io.jsonwebtoken.Jws; |
||||
|
import io.jsonwebtoken.Jwts; |
||||
|
import io.jsonwebtoken.MalformedJwtException; |
||||
|
import io.jsonwebtoken.SignatureException; |
||||
|
import io.jsonwebtoken.UnsupportedJwtException; |
||||
|
import io.micrometer.core.instrument.util.StringUtils; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
|
||||
|
import java.util.Date; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
@Component |
||||
|
@Slf4j |
||||
|
public class OAuth2AppTokenFactory { |
||||
|
|
||||
|
private static final String CALLBACK_URL_SCHEME = "callbackUrlScheme"; |
||||
|
|
||||
|
private static final long MAX_EXPIRATION_TIME_DIFF_MS = TimeUnit.MINUTES.toMillis(5); |
||||
|
|
||||
|
public String validateTokenAndGetCallbackUrlScheme(String appPackage, String appToken, String appSecret) { |
||||
|
Jws<Claims> jwsClaims; |
||||
|
try { |
||||
|
jwsClaims = Jwts.parser().setSigningKey(appSecret).parseClaimsJws(appToken); |
||||
|
} |
||||
|
catch (UnsupportedJwtException | MalformedJwtException | IllegalArgumentException | SignatureException ex) { |
||||
|
throw new IllegalArgumentException("Invalid Application token: ", ex); |
||||
|
} catch (ExpiredJwtException expiredEx) { |
||||
|
throw new IllegalArgumentException("Application token expired", expiredEx); |
||||
|
} |
||||
|
Claims claims = jwsClaims.getBody(); |
||||
|
Date expiration = claims.getExpiration(); |
||||
|
if (expiration == null) { |
||||
|
throw new IllegalArgumentException("Application token must have expiration date"); |
||||
|
} |
||||
|
long timeDiff = expiration.getTime() - System.currentTimeMillis(); |
||||
|
if (timeDiff > MAX_EXPIRATION_TIME_DIFF_MS) { |
||||
|
throw new IllegalArgumentException("Application token expiration time can't be longer than 5 minutes"); |
||||
|
} |
||||
|
if (!claims.getIssuer().equals(appPackage)) { |
||||
|
throw new IllegalArgumentException("Application token issuer doesn't match application package"); |
||||
|
} |
||||
|
String callbackUrlScheme = claims.get(CALLBACK_URL_SCHEME, String.class); |
||||
|
if (StringUtils.isEmpty(callbackUrlScheme)) { |
||||
|
throw new IllegalArgumentException("Application token doesn't have callbackUrlScheme"); |
||||
|
} |
||||
|
return callbackUrlScheme; |
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue