committed by
GitHub
1228 changed files with 46819 additions and 14265 deletions
@ -1,162 +0,0 @@ |
|||
{ |
|||
"ruleChain": { |
|||
"additionalInfo": null, |
|||
"name": "Edge Root Rule Chain", |
|||
"type": "EDGE", |
|||
"firstRuleNodeId": null, |
|||
"root": true, |
|||
"debugMode": false, |
|||
"configuration": null |
|||
}, |
|||
"metadata": { |
|||
"firstNodeIndex": 0, |
|||
"nodes": [ |
|||
{ |
|||
"additionalInfo": { |
|||
"description": "Process incoming messages from devices with the alarm rules defined in the device profile. Dispatch all incoming messages with \"Success\" relation type.", |
|||
"layoutX": 203, |
|||
"layoutY": 259 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.profile.TbDeviceProfileNode", |
|||
"name": "Device Profile Node", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"persistAlarmRulesState": false, |
|||
"fetchAlarmRulesStateOnStart": false |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 823, |
|||
"layoutY": 157 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode", |
|||
"name": "Save Timeseries", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"defaultTTL": 0 |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 824, |
|||
"layoutY": 52 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode", |
|||
"name": "Save Client Attributes", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"scope": "CLIENT_SCOPE" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 347, |
|||
"layoutY": 149 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode", |
|||
"name": "Message Type Switch", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"version": 0 |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 825, |
|||
"layoutY": 266 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbLogNode", |
|||
"name": "Log RPC from Device", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 824, |
|||
"layoutY": 378 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbLogNode", |
|||
"name": "Log Other", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 824, |
|||
"layoutY": 466 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.rpc.TbSendRPCRequestNode", |
|||
"name": "RPC Call Request", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"timeoutInSeconds": 60 |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 1134, |
|||
"layoutY": 132 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.edge.TbMsgPushToCloudNode", |
|||
"name": "Push to cloud", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"version": 0 |
|||
} |
|||
} |
|||
], |
|||
"connections": [ |
|||
{ |
|||
"fromIndex": 0, |
|||
"toIndex": 3, |
|||
"type": "Success" |
|||
}, |
|||
{ |
|||
"fromIndex": 1, |
|||
"toIndex": 7, |
|||
"type": "Success" |
|||
}, |
|||
{ |
|||
"fromIndex": 2, |
|||
"toIndex": 7, |
|||
"type": "Success" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 6, |
|||
"type": "RPC Request to Device" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 5, |
|||
"type": "Other" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 2, |
|||
"type": "Post attributes" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 1, |
|||
"type": "Post telemetry" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 4, |
|||
"type": "RPC Request from Device" |
|||
}, |
|||
{ |
|||
"fromIndex": 4, |
|||
"toIndex": 7, |
|||
"type": "Success" |
|||
} |
|||
], |
|||
"ruleChainConnections": null |
|||
} |
|||
} |
|||
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
File diff suppressed because one or more lines are too long
File diff suppressed because it is too large
@ -0,0 +1,70 @@ |
|||
/** |
|||
* 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 io.swagger.annotations.ApiOperation; |
|||
import io.swagger.annotations.ApiParam; |
|||
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.RequestBody; |
|||
import org.springframework.web.bind.annotation.RequestMapping; |
|||
import org.springframework.web.bind.annotation.RequestMethod; |
|||
import org.springframework.web.bind.annotation.ResponseBody; |
|||
import org.springframework.web.bind.annotation.RestController; |
|||
import org.springframework.web.context.request.async.DeferredResult; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID_PARAM_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH; |
|||
|
|||
@RestController |
|||
@TbCoreComponent |
|||
@RequestMapping(TbUrlConstants.RPC_V1_URL_PREFIX) |
|||
@Slf4j |
|||
public class RpcV1Controller extends AbstractRpcController { |
|||
|
|||
@ApiOperation(value = "Send one-way RPC request (handleOneWayDeviceRPCRequest)", notes = "Deprecated. See 'Rpc V 2 Controller' instead." + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH) |
|||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/oneway/{deviceId}", method = RequestMethod.POST) |
|||
@ResponseBody |
|||
public DeferredResult<ResponseEntity> handleOneWayDeviceRPCRequest( |
|||
@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable("deviceId") String deviceIdStr, |
|||
@ApiParam(value = "A JSON value representing the RPC request.") |
|||
@RequestBody String requestBody) throws ThingsboardException { |
|||
return handleDeviceRPCRequest(true, new DeviceId(UUID.fromString(deviceIdStr)), requestBody, HttpStatus.REQUEST_TIMEOUT, HttpStatus.CONFLICT); |
|||
} |
|||
|
|||
@ApiOperation(value = "Send two-way RPC request (handleTwoWayDeviceRPCRequest)", notes = "Deprecated. See 'Rpc V 2 Controller' instead." + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH) |
|||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/twoway/{deviceId}", method = RequestMethod.POST) |
|||
@ResponseBody |
|||
public DeferredResult<ResponseEntity> handleTwoWayDeviceRPCRequest( |
|||
@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable("deviceId") String deviceIdStr, |
|||
@ApiParam(value = "A JSON value representing the RPC request.") |
|||
@RequestBody String requestBody) throws ThingsboardException { |
|||
return handleDeviceRPCRequest(false, new DeviceId(UUID.fromString(deviceIdStr)), requestBody, HttpStatus.REQUEST_TIMEOUT, HttpStatus.CONFLICT); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,246 @@ |
|||
/** |
|||
* 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 com.google.common.util.concurrent.FutureCallback; |
|||
import io.swagger.annotations.ApiOperation; |
|||
import io.swagger.annotations.ApiParam; |
|||
import io.swagger.annotations.ApiResponse; |
|||
import io.swagger.annotations.ApiResponses; |
|||
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.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.context.request.async.DeferredResult; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
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.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.rpc.RemoveRpcActorMsg; |
|||
import org.thingsboard.server.service.security.permission.Operation; |
|||
import org.thingsboard.server.service.telemetry.exception.ToErrorResponseEntity; |
|||
|
|||
import javax.annotation.Nullable; |
|||
import java.util.UUID; |
|||
|
|||
import static org.thingsboard.server.common.data.DataConstants.RPC_DELETED; |
|||
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID; |
|||
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID_PARAM_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.MARKDOWN_CODE_BLOCK_END; |
|||
import static org.thingsboard.server.controller.ControllerConstants.MARKDOWN_CODE_BLOCK_START; |
|||
import static org.thingsboard.server.controller.ControllerConstants.PAGE_NUMBER_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.PAGE_SIZE_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.RPC_ID; |
|||
import static org.thingsboard.server.controller.ControllerConstants.RPC_ID_PARAM_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.RPC_SORT_PROPERTY_ALLOWABLE_VALUES; |
|||
import static org.thingsboard.server.controller.ControllerConstants.RPC_STATUS_ALLOWABLE_VALUES; |
|||
import static org.thingsboard.server.controller.ControllerConstants.RPC_TEXT_SEARCH_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.SORT_ORDER_ALLOWABLE_VALUES; |
|||
import static org.thingsboard.server.controller.ControllerConstants.SORT_ORDER_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.SORT_PROPERTY_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHORITY_PARAGRAPH; |
|||
import static org.thingsboard.server.controller.ControllerConstants.TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH; |
|||
|
|||
@RestController |
|||
@TbCoreComponent |
|||
@RequestMapping(TbUrlConstants.RPC_V2_URL_PREFIX) |
|||
@Slf4j |
|||
public class RpcV2Controller extends AbstractRpcController { |
|||
|
|||
private static final String RPC_REQUEST_DESCRIPTION = "Sends the one-way remote-procedure call (RPC) request to device. " + |
|||
"The RPC call is A JSON that contains the method name ('method'), parameters ('params') and multiple optional fields. " + |
|||
"See example below. We will review the properties of the RPC call one-by-one below. " + |
|||
"\n\n" + MARKDOWN_CODE_BLOCK_START + |
|||
"{\n" + |
|||
" \"method\": \"setGpio\",\n" + |
|||
" \"params\": {\n" + |
|||
" \"pin\": 7,\n" + |
|||
" \"value\": 1\n" + |
|||
" },\n" + |
|||
" \"persistent\": false,\n" + |
|||
" \"timeout\": 5000\n" + |
|||
"}" + |
|||
MARKDOWN_CODE_BLOCK_END + |
|||
"\n\n### Server-side RPC structure\n" + |
|||
"\n" + |
|||
"The body of server-side RPC request consists of multiple fields:\n" + |
|||
"\n" + |
|||
"* **method** - mandatory, name of the method to distinct the RPC calls.\n" + |
|||
" For example, \"getCurrentTime\" or \"getWeatherForecast\". The value of the parameter is a string.\n" + |
|||
"* **params** - mandatory, parameters used for processing of the request. The value is a JSON. Leave empty JSON \"{}\" if no parameters needed.\n" + |
|||
"* **timeout** - optional, value of the processing timeout in milliseconds. The default value is 10000 (10 seconds). The minimum value is 5000 (5 seconds).\n" + |
|||
"* **expirationTime** - optional, value of the epoch time (in milliseconds, UTC timezone). Overrides **timeout** if present.\n" + |
|||
"* **persistent** - optional, indicates persistent RPC. The default value is \"false\".\n" + |
|||
"* **retries** - optional, defines how many times persistent RPC will be re-sent in case of failures on the network and/or device side.\n" + |
|||
"* **additionalInfo** - optional, defines metadata for the persistent RPC that will be added to the persistent RPC events."; |
|||
|
|||
private static final String ONE_WAY_RPC_RESULT = "\n\n### RPC Result\n" + |
|||
"In case of persistent RPC, the result of this call is 'rpcId' UUID. In case of lightweight RPC, " + |
|||
"the result of this call is either 200 OK if the message was sent to device, or 504 Gateway Timeout if device is offline."; |
|||
|
|||
private static final String TWO_WAY_RPC_RESULT = "\n\n### RPC Result\n" + |
|||
"In case of persistent RPC, the result of this call is 'rpcId' UUID. In case of lightweight RPC, " + |
|||
"the result of this call is the response from device, or 504 Gateway Timeout if device is offline."; |
|||
|
|||
private static final String ONE_WAY_RPC_REQUEST_DESCRIPTION = "Sends the one-way remote-procedure call (RPC) request to device. " + RPC_REQUEST_DESCRIPTION + ONE_WAY_RPC_RESULT + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH; |
|||
|
|||
private static final String TWO_WAY_RPC_REQUEST_DESCRIPTION = "Sends the two-way remote-procedure call (RPC) request to device. " + RPC_REQUEST_DESCRIPTION + TWO_WAY_RPC_RESULT + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH; |
|||
|
|||
@ApiOperation(value = "Send one-way RPC request", notes = ONE_WAY_RPC_REQUEST_DESCRIPTION) |
|||
@ApiResponses(value = { |
|||
@ApiResponse(code = 200, message = "Persistent RPC request was saved to the database or lightweight RPC request was sent to the device."), |
|||
@ApiResponse(code = 400, message = "Invalid structure of the request."), |
|||
@ApiResponse(code = 401, message = "User is not authorized to send the RPC request. Most likely, User belongs to different Customer or Tenant."), |
|||
@ApiResponse(code = 504, message = "Timeout to process the RPC call. Most likely, device is offline."), |
|||
}) |
|||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/oneway/{deviceId}", method = RequestMethod.POST) |
|||
@ResponseBody |
|||
public DeferredResult<ResponseEntity> handleOneWayDeviceRPCRequest( |
|||
@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable("deviceId") String deviceIdStr, |
|||
@ApiParam(value = "A JSON value representing the RPC request.") |
|||
@RequestBody String requestBody) throws ThingsboardException { |
|||
return handleDeviceRPCRequest(true, new DeviceId(UUID.fromString(deviceIdStr)), requestBody, HttpStatus.GATEWAY_TIMEOUT, HttpStatus.GATEWAY_TIMEOUT); |
|||
} |
|||
|
|||
@ApiOperation(value = "Send two-way RPC request", notes = TWO_WAY_RPC_REQUEST_DESCRIPTION) |
|||
@ApiResponses(value = { |
|||
@ApiResponse(code = 200, message = "Persistent RPC request was saved to the database or lightweight RPC response received."), |
|||
@ApiResponse(code = 400, message = "Invalid structure of the request."), |
|||
@ApiResponse(code = 401, message = "User is not authorized to send the RPC request. Most likely, User belongs to different Customer or Tenant."), |
|||
@ApiResponse(code = 504, message = "Timeout to process the RPC call. Most likely, device is offline."), |
|||
}) |
|||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/twoway/{deviceId}", method = RequestMethod.POST) |
|||
@ResponseBody |
|||
public DeferredResult<ResponseEntity> handleTwoWayDeviceRPCRequest( |
|||
@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable(DEVICE_ID) String deviceIdStr, |
|||
@ApiParam(value = "A JSON value representing the RPC request.") |
|||
@RequestBody String requestBody) throws ThingsboardException { |
|||
return handleDeviceRPCRequest(false, new DeviceId(UUID.fromString(deviceIdStr)), requestBody, HttpStatus.GATEWAY_TIMEOUT, HttpStatus.GATEWAY_TIMEOUT); |
|||
} |
|||
|
|||
@ApiOperation(value = "Get persistent RPC request", notes = "Get information about the status of the RPC call." + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH) |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/persistent/{rpcId}", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public Rpc getPersistedRpc( |
|||
@ApiParam(value = RPC_ID_PARAM_DESCRIPTION, required = true) |
|||
@PathVariable(RPC_ID) String strRpc) throws ThingsboardException { |
|||
checkParameter("RpcId", strRpc); |
|||
try { |
|||
RpcId rpcId = new RpcId(UUID.fromString(strRpc)); |
|||
return checkRpcId(rpcId, Operation.READ); |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
|
|||
@ApiOperation(value = "Get persistent RPC requests", notes = "Allows to query RPC calls for specific device using pagination." + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH) |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/persistent/device/{deviceId}", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public DeferredResult<ResponseEntity> getPersistedRpcByDevice( |
|||
@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION, required = true) |
|||
@PathVariable(DEVICE_ID) String strDeviceId, |
|||
@ApiParam(value = PAGE_SIZE_DESCRIPTION, required = true) |
|||
@RequestParam int pageSize, |
|||
@ApiParam(value = PAGE_NUMBER_DESCRIPTION, required = true) |
|||
@RequestParam int page, |
|||
@ApiParam(value = "Status of the RPC", required = true, allowableValues = RPC_STATUS_ALLOWABLE_VALUES) |
|||
@RequestParam RpcStatus rpcStatus, |
|||
@ApiParam(value = RPC_TEXT_SEARCH_DESCRIPTION) |
|||
@RequestParam(required = false) String textSearch, |
|||
@ApiParam(value = SORT_PROPERTY_DESCRIPTION, allowableValues = RPC_SORT_PROPERTY_ALLOWABLE_VALUES) |
|||
@RequestParam(required = false) String sortProperty, |
|||
@ApiParam(value = SORT_ORDER_DESCRIPTION, allowableValues = SORT_ORDER_ALLOWABLE_VALUES) |
|||
@RequestParam(required = false) String sortOrder) throws ThingsboardException { |
|||
checkParameter("DeviceId", strDeviceId); |
|||
try { |
|||
TenantId tenantId = getCurrentUser().getTenantId(); |
|||
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); |
|||
DeviceId deviceId = new DeviceId(UUID.fromString(strDeviceId)); |
|||
final DeferredResult<ResponseEntity> response = new DeferredResult<>(); |
|||
accessValidator.validate(getCurrentUser(), Operation.RPC_CALL, deviceId, new HttpValidationCallback(response, new FutureCallback<>() { |
|||
@Override |
|||
public void onSuccess(@Nullable DeferredResult<ResponseEntity> result) { |
|||
PageData<Rpc> rpcCalls = rpcService.findAllByDeviceIdAndStatus(tenantId, deviceId, rpcStatus, pageLink); |
|||
response.setResult(new ResponseEntity<>(rpcCalls, HttpStatus.OK)); |
|||
} |
|||
|
|||
@Override |
|||
public void onFailure(Throwable e) { |
|||
ResponseEntity entity; |
|||
if (e instanceof ToErrorResponseEntity) { |
|||
entity = ((ToErrorResponseEntity) e).toErrorResponseEntity(); |
|||
} else { |
|||
entity = new ResponseEntity(HttpStatus.UNAUTHORIZED); |
|||
} |
|||
response.setResult(entity); |
|||
} |
|||
})); |
|||
return response; |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
|
|||
@ApiOperation(value = "Delete persistent RPC", notes = "Deletes the persistent RPC request." + TENANT_AUTHORITY_PARAGRAPH) |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/persistent/{rpcId}", method = RequestMethod.DELETE) |
|||
@ResponseBody |
|||
public void deleteResource( |
|||
@ApiParam(value = RPC_ID_PARAM_DESCRIPTION, required = true) |
|||
@PathVariable(RPC_ID) String strRpc) throws ThingsboardException { |
|||
checkParameter("RpcId", strRpc); |
|||
try { |
|||
RpcId rpcId = new RpcId(UUID.fromString(strRpc)); |
|||
Rpc rpc = checkRpcId(rpcId, Operation.DELETE); |
|||
|
|||
if (rpc != null) { |
|||
if (rpc.getStatus().equals(RpcStatus.QUEUED)) { |
|||
RemoveRpcActorMsg removeMsg = new RemoveRpcActorMsg(getTenantId(), rpc.getDeviceId(), rpc.getUuidId()); |
|||
log.trace("[{}] Forwarding msg {} to queue actor!", rpc.getDeviceId(), rpc); |
|||
tbClusterService.pushMsgToCore(removeMsg, null); |
|||
} |
|||
|
|||
rpcService.deleteRpc(getTenantId(), rpcId); |
|||
|
|||
TbMsg msg = TbMsg.newMsg(RPC_DELETED, rpc.getDeviceId(), TbMsgMetaData.EMPTY, JacksonUtil.toString(rpc)); |
|||
tbClusterService.pushMsgToRuleEngine(getTenantId(), rpc.getDeviceId(), msg, null); |
|||
} |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
} |
|||
@ -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.controller; |
|||
|
|||
import io.swagger.annotations.ApiOperation; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.security.access.prepost.PreAuthorize; |
|||
import org.springframework.web.bind.annotation.RequestMapping; |
|||
import org.springframework.web.bind.annotation.RequestMethod; |
|||
import org.springframework.web.bind.annotation.ResponseBody; |
|||
import org.springframework.web.bind.annotation.RestController; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
@RestController |
|||
@TbCoreComponent |
|||
@RequestMapping("/api") |
|||
public class UiSettingsController extends BaseController { |
|||
|
|||
@Value("${ui.help.base-url}") |
|||
private String helpBaseUrl; |
|||
|
|||
@ApiOperation(value = "Get UI help base url (getHelpBaseUrl)", |
|||
notes = "Get UI help base url used to fetch help assets. " + |
|||
"The actual value of the base url is configurable in the system configuration file.") |
|||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/uiSettings/helpBaseUrl", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public String getHelpBaseUrl() throws ThingsboardException { |
|||
return helpBaseUrl; |
|||
} |
|||
|
|||
} |
|||
@ -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<Asset> { |
|||
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<Asset> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user) { |
|||
ImportedEntityInfo<Asset> 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<BulkImportColumnType, String> 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); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,276 @@ |
|||
/** |
|||
* 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; |
|||
import java.util.concurrent.locks.Lock; |
|||
import java.util.concurrent.locks.ReentrantLock; |
|||
|
|||
@Service |
|||
@TbCoreComponent |
|||
public class DeviceBulkImportService extends AbstractBulkImportService<Device> { |
|||
protected final DeviceService deviceService; |
|||
protected final DeviceCredentialsService deviceCredentialsService; |
|||
protected final DeviceProfileService deviceProfileService; |
|||
|
|||
private final Lock findOrCreateDeviceProfileLock = new ReentrantLock(); |
|||
|
|||
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<Device> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user) { |
|||
ImportedEntityInfo<Device> 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()); |
|||
} |
|||
|
|||
DeviceProfile deviceProfile; |
|||
if (deviceCredentials.getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) { |
|||
deviceProfile = setUpLwM2mDeviceProfile(user.getTenantId(), device); |
|||
} else if (StringUtils.isNotEmpty(device.getType())) { |
|||
deviceProfile = deviceProfileService.findOrCreateDeviceProfile(user.getTenantId(), device.getType()); |
|||
} else { |
|||
deviceProfile = deviceProfileService.findDefaultDeviceProfile(user.getTenantId()); |
|||
} |
|||
device.setDeviceProfileId(deviceProfile.getId()); |
|||
|
|||
device = deviceService.saveDeviceWithCredentials(device, deviceCredentials); |
|||
|
|||
importedEntityInfo.setEntity(device); |
|||
return importedEntityInfo; |
|||
} |
|||
|
|||
private void setDeviceFields(Device device, Map<BulkImportColumnType, String> 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<BulkImportColumnType, String> 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<BulkImportColumnType, String> fields, DeviceCredentials credentials) { |
|||
credentials.setCredentialsId(Optional.ofNullable(fields.get(BulkImportColumnType.ACCESS_TOKEN)) |
|||
.orElseGet(() -> RandomStringUtils.randomAlphanumeric(20))); |
|||
} |
|||
|
|||
private void setUpBasicMqttCredentials(Map<BulkImportColumnType, String> 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<BulkImportColumnType, String> fields, DeviceCredentials credentials) { |
|||
credentials.setCredentialsValue(fields.get(BulkImportColumnType.X509)); |
|||
} |
|||
|
|||
private void setUpLwm2mCredentials(Map<BulkImportColumnType, String> 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 DeviceProfile 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); |
|||
} |
|||
} else { |
|||
findOrCreateDeviceProfileLock.lock(); |
|||
try { |
|||
deviceProfile = deviceProfileService.findDeviceProfileByName(tenantId, device.getType()); |
|||
if (deviceProfile == null) { |
|||
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); |
|||
} |
|||
} finally { |
|||
findOrCreateDeviceProfileLock.unlock(); |
|||
} |
|||
} |
|||
return deviceProfile; |
|||
} |
|||
|
|||
private void setValues(ObjectNode objectNode, Map<BulkImportColumnType, String> data, Collection<BulkImportColumnType> 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)); |
|||
} |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -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.edge; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.http.HttpHost; |
|||
import org.apache.http.conn.ssl.DefaultHostnameVerifier; |
|||
import org.apache.http.impl.client.CloseableHttpClient; |
|||
import org.apache.http.impl.client.HttpClients; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.http.ResponseEntity; |
|||
import org.springframework.http.client.HttpComponentsClientHttpRequestFactory; |
|||
import org.springframework.http.client.SimpleClientHttpRequestFactory; |
|||
import org.springframework.stereotype.Service; |
|||
import org.springframework.web.client.RestTemplate; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import java.net.InetSocketAddress; |
|||
import java.net.Proxy; |
|||
import java.util.HashMap; |
|||
import java.util.Map; |
|||
|
|||
import static org.apache.commons.lang3.StringUtils.isNotEmpty; |
|||
|
|||
@Service |
|||
@TbCoreComponent |
|||
@Slf4j |
|||
public class DefaultEdgeLicenseService implements EdgeLicenseService { |
|||
|
|||
private RestTemplate restTemplate; |
|||
|
|||
private static final String EDGE_LICENSE_SERVER_ENDPOINT = "https://license.thingsboard.io"; |
|||
|
|||
@Value("${edges.enabled:false}") |
|||
private boolean edgesEnabled; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
if (edgesEnabled) { |
|||
initRestTemplate(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public ResponseEntity<JsonNode> checkInstance(JsonNode request) { |
|||
return this.restTemplate.postForEntity(EDGE_LICENSE_SERVER_ENDPOINT + "/api/license/checkInstance", request, JsonNode.class); |
|||
} |
|||
|
|||
@Override |
|||
public ResponseEntity<JsonNode> activateInstance(String edgeLicenseSecret, String releaseDate) { |
|||
Map<String, String> params = new HashMap<>(); |
|||
params.put("licenseSecret", edgeLicenseSecret); |
|||
params.put("releaseDate", releaseDate); |
|||
return this.restTemplate.postForEntity(EDGE_LICENSE_SERVER_ENDPOINT + "/api/license/activateInstance?licenseSecret={licenseSecret}&releaseDate={releaseDate}", null, JsonNode.class, params); |
|||
} |
|||
|
|||
private void initRestTemplate() { |
|||
boolean jdkHttpClientEnabled = isNotEmpty(System.getProperty("tb.proxy.jdk")) && System.getProperty("tb.proxy.jdk").equalsIgnoreCase("true"); |
|||
boolean systemProxyEnabled = isNotEmpty(System.getProperty("tb.proxy.system")) && System.getProperty("tb.proxy.system").equalsIgnoreCase("true"); |
|||
boolean proxyEnabled = isNotEmpty(System.getProperty("tb.proxy.host")) && isNotEmpty(System.getProperty("tb.proxy.port")); |
|||
if (jdkHttpClientEnabled) { |
|||
log.warn("Going to use plain JDK Http Client!"); |
|||
SimpleClientHttpRequestFactory factory = new SimpleClientHttpRequestFactory(); |
|||
if (proxyEnabled) { |
|||
log.warn("Going to use Proxy Server: [{}:{}]", System.getProperty("tb.proxy.host"), System.getProperty("tb.proxy.port")); |
|||
factory.setProxy(new Proxy(Proxy.Type.HTTP, InetSocketAddress.createUnresolved(System.getProperty("tb.proxy.host"), Integer.parseInt(System.getProperty("tb.proxy.port"))))); |
|||
} |
|||
|
|||
this.restTemplate = new RestTemplate(new SimpleClientHttpRequestFactory()); |
|||
} else { |
|||
CloseableHttpClient httpClient; |
|||
HttpComponentsClientHttpRequestFactory requestFactory; |
|||
if (systemProxyEnabled) { |
|||
log.warn("Going to use System Proxy Server!"); |
|||
httpClient = HttpClients.createSystem(); |
|||
requestFactory = new HttpComponentsClientHttpRequestFactory(); |
|||
requestFactory.setHttpClient(httpClient); |
|||
this.restTemplate = new RestTemplate(requestFactory); |
|||
} else if (proxyEnabled) { |
|||
log.warn("Going to use Proxy Server: [{}:{}]", System.getProperty("tb.proxy.host"), System.getProperty("tb.proxy.port")); |
|||
httpClient = HttpClients.custom().setSSLHostnameVerifier(new DefaultHostnameVerifier()).setProxy(new HttpHost(System.getProperty("tb.proxy.host"), Integer.parseInt(System.getProperty("tb.proxy.port")), "https")).build(); |
|||
requestFactory = new HttpComponentsClientHttpRequestFactory(); |
|||
requestFactory.setHttpClient(httpClient); |
|||
this.restTemplate = new RestTemplate(requestFactory); |
|||
} else { |
|||
httpClient = HttpClients.custom().setSSLHostnameVerifier(new DefaultHostnameVerifier()).build(); |
|||
requestFactory = new HttpComponentsClientHttpRequestFactory(); |
|||
requestFactory.setHttpClient(httpClient); |
|||
this.restTemplate = new RestTemplate(requestFactory); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
|
|||
|
|||
@ -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<Edge> { |
|||
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<Edge> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user) { |
|||
ImportedEntityInfo<Edge> 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<BulkImportColumnType, String> 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); |
|||
} |
|||
|
|||
} |
|||
@ -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.edge; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import org.springframework.http.ResponseEntity; |
|||
|
|||
public interface EdgeLicenseService { |
|||
|
|||
ResponseEntity<JsonNode> checkInstance(JsonNode request); |
|||
|
|||
ResponseEntity<JsonNode> activateInstance(String licenseSecret, String releaseDate); |
|||
} |
|||
@ -0,0 +1,48 @@ |
|||
/** |
|||
* 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.rpc; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.id.EdgeId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
|
|||
public final class EdgeEventUtils { |
|||
|
|||
private EdgeEventUtils() { |
|||
} |
|||
|
|||
public static EdgeEvent constructEdgeEvent(TenantId tenantId, |
|||
EdgeId edgeId, |
|||
EdgeEventType type, |
|||
EdgeEventActionType action, |
|||
EntityId entityId, |
|||
JsonNode body) { |
|||
EdgeEvent edgeEvent = new EdgeEvent(); |
|||
edgeEvent.setTenantId(tenantId); |
|||
edgeEvent.setEdgeId(edgeId); |
|||
edgeEvent.setType(type); |
|||
edgeEvent.setAction(action); |
|||
if (entityId != null) { |
|||
edgeEvent.setEntityId(entityId.getId()); |
|||
} |
|||
edgeEvent.setBody(body); |
|||
return edgeEvent; |
|||
} |
|||
} |
|||
File diff suppressed because it is too large
@ -0,0 +1,32 @@ |
|||
/** |
|||
* 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.rpc; |
|||
|
|||
import com.google.common.util.concurrent.SettableFuture; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|||
|
|||
import java.util.LinkedHashMap; |
|||
import java.util.Map; |
|||
import java.util.concurrent.ScheduledFuture; |
|||
|
|||
@Data |
|||
public class EdgeSessionState { |
|||
|
|||
private final Map<Integer, DownlinkMsg> pendingMsgsMap = new LinkedHashMap<>(); |
|||
private SettableFuture<Void> sendDownlinkMsgsFuture; |
|||
private ScheduledFuture<?> scheduledSendDownlinkTask; |
|||
} |
|||
@ -0,0 +1,74 @@ |
|||
/** |
|||
* 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.rpc; |
|||
|
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.service.edge.EdgeContextComponent; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.AssetsEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.CustomerEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.CustomerUsersEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher; |
|||
import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetsBundlesEdgeEventFetcher; |
|||
|
|||
import java.util.LinkedList; |
|||
import java.util.List; |
|||
import java.util.NoSuchElementException; |
|||
|
|||
public class EdgeSyncCursor { |
|||
|
|||
List<EdgeEventFetcher> fetchers = new LinkedList<>(); |
|||
|
|||
int currentIdx = 0; |
|||
|
|||
public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge) { |
|||
fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); |
|||
fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); |
|||
fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); |
|||
fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); |
|||
fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); |
|||
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { |
|||
fetchers.add(new CustomerEdgeEventFetcher()); |
|||
fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId())); |
|||
} |
|||
fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), ctx.getFreemarkerConfig())); |
|||
fetchers.add(new AssetsEdgeEventFetcher(ctx.getAssetService())); |
|||
fetchers.add(new DashboardsEdgeEventFetcher(ctx.getDashboardService())); |
|||
} |
|||
|
|||
public boolean hasNext() { |
|||
return fetchers.size() > currentIdx; |
|||
} |
|||
|
|||
public EdgeEventFetcher getNext() { |
|||
if (!hasNext()) { |
|||
throw new NoSuchElementException(); |
|||
} |
|||
EdgeEventFetcher edgeEventFetcher = fetchers.get(currentIdx); |
|||
currentIdx++; |
|||
return edgeEventFetcher; |
|||
} |
|||
|
|||
public int getCurrentIdx() { |
|||
return currentIdx; |
|||
} |
|||
} |
|||
@ -0,0 +1,161 @@ |
|||
/** |
|||
* 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.rpc.fetch; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import com.fasterxml.jackson.databind.ObjectMapper; |
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import freemarker.template.Configuration; |
|||
import freemarker.template.Template; |
|||
import lombok.AllArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.apache.commons.lang3.text.WordUtils; |
|||
import org.thingsboard.server.common.data.AdminSettings; |
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.id.AdminSettingsId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.dao.settings.AdminSettingsService; |
|||
import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Arrays; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.regex.Matcher; |
|||
import java.util.regex.Pattern; |
|||
|
|||
@AllArgsConstructor |
|||
@Slf4j |
|||
public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher { |
|||
|
|||
private static final ObjectMapper mapper = new ObjectMapper(); |
|||
|
|||
private final AdminSettingsService adminSettingsService; |
|||
private final Configuration freemarkerConfig; |
|||
|
|||
private static Pattern startPattern = Pattern.compile("<div class=\"content\".*?>"); |
|||
private static Pattern endPattern = Pattern.compile("<div class=\"footer\".*?>"); |
|||
|
|||
private static List<String> templatesNames = Arrays.asList( |
|||
"account.activated.ftl", |
|||
"account.lockout.ftl", |
|||
"activation.ftl", |
|||
"password.was.reset.ftl", |
|||
"reset.password.ftl", |
|||
"test.ftl"); |
|||
|
|||
// TODO: fix format of next templates
|
|||
// "state.disabled.ftl",
|
|||
// "state.enabled.ftl",
|
|||
// "state.warning.ftl",
|
|||
|
|||
@Override |
|||
public PageLink getPageLink(int pageSize) { |
|||
return null; |
|||
} |
|||
|
|||
@Override |
|||
public PageData<EdgeEvent> fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) throws Exception { |
|||
List<EdgeEvent> result = new ArrayList<>(); |
|||
|
|||
AdminSettings systemMailSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "mail"); |
|||
result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, |
|||
EdgeEventActionType.UPDATED, null, mapper.valueToTree(systemMailSettings))); |
|||
|
|||
AdminSettings tenantMailSettings = convertToTenantAdminSettings(systemMailSettings.getKey(), (ObjectNode) systemMailSettings.getJsonValue()); |
|||
result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, |
|||
EdgeEventActionType.UPDATED, null, mapper.valueToTree(tenantMailSettings))); |
|||
|
|||
AdminSettings systemMailTemplates = loadMailTemplates(); |
|||
result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, |
|||
EdgeEventActionType.UPDATED, null, mapper.valueToTree(systemMailTemplates))); |
|||
|
|||
AdminSettings tenantMailTemplates = convertToTenantAdminSettings(systemMailTemplates.getKey(), (ObjectNode) systemMailTemplates.getJsonValue()); |
|||
result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, |
|||
EdgeEventActionType.UPDATED, null, mapper.valueToTree(tenantMailTemplates))); |
|||
|
|||
// @voba - returns PageData object to be in sync with other fetchers
|
|||
return new PageData<>(result, 1, result.size(), false); |
|||
} |
|||
|
|||
private AdminSettings loadMailTemplates() throws Exception { |
|||
Map<String, Object> mailTemplates = new HashMap<>(); |
|||
for (String templatesName : templatesNames) { |
|||
Template template = freemarkerConfig.getTemplate(templatesName); |
|||
if (template != null) { |
|||
String name = validateName(template.getName()); |
|||
Map<String, String> mailTemplate = getMailTemplateFromFile(template.toString()); |
|||
if (mailTemplate != null) { |
|||
mailTemplates.put(name, mailTemplate); |
|||
} else { |
|||
log.error("Can't load mail template from file {}", template.getName()); |
|||
} |
|||
} |
|||
} |
|||
AdminSettings adminSettings = new AdminSettings(); |
|||
adminSettings.setId(new AdminSettingsId(Uuids.timeBased())); |
|||
adminSettings.setKey("mailTemplates"); |
|||
adminSettings.setJsonValue(mapper.convertValue(mailTemplates, JsonNode.class)); |
|||
return adminSettings; |
|||
} |
|||
|
|||
private Map<String, String> getMailTemplateFromFile(String stringTemplate) { |
|||
Map<String, String> mailTemplate = new HashMap<>(); |
|||
Matcher start = startPattern.matcher(stringTemplate); |
|||
Matcher end = endPattern.matcher(stringTemplate); |
|||
if (start.find() && end.find()) { |
|||
String body = StringUtils.substringBetween(stringTemplate, start.group(), end.group()).replaceAll("\t", ""); |
|||
String subject = StringUtils.substringBetween(body, "<h2>", "</h2>"); |
|||
mailTemplate.put("subject", subject); |
|||
mailTemplate.put("body", body); |
|||
} else { |
|||
return null; |
|||
} |
|||
return mailTemplate; |
|||
} |
|||
|
|||
private String validateName(String name) throws Exception { |
|||
StringBuilder nameBuilder = new StringBuilder(); |
|||
name = name.replace(".ftl", ""); |
|||
String[] nameParts = name.split("\\."); |
|||
if (nameParts.length >= 1) { |
|||
nameBuilder.append(nameParts[0]); |
|||
for (int i = 1; i < nameParts.length; i++) { |
|||
String word = WordUtils.capitalize(nameParts[i]); |
|||
nameBuilder.append(word); |
|||
} |
|||
return nameBuilder.toString(); |
|||
} else { |
|||
throw new Exception("Error during filename validation"); |
|||
} |
|||
} |
|||
|
|||
private AdminSettings convertToTenantAdminSettings(String key, ObjectNode jsonValue) { |
|||
AdminSettings tenantMailSettings = new AdminSettings(); |
|||
jsonValue.put("useSystemMailSettings", true); |
|||
tenantMailSettings.setJsonValue(jsonValue); |
|||
tenantMailSettings.setKey(key); |
|||
return tenantMailSettings; |
|||
} |
|||
} |
|||
@ -0,0 +1,47 @@ |
|||
/** |
|||
* 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.rpc.fetch; |
|||
|
|||
import lombok.AllArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.dao.asset.AssetService; |
|||
import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; |
|||
|
|||
@AllArgsConstructor |
|||
@Slf4j |
|||
public class AssetsEdgeEventFetcher extends BasePageableEdgeEventFetcher<Asset> { |
|||
|
|||
private final AssetService assetService; |
|||
|
|||
@Override |
|||
PageData<Asset> fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { |
|||
return assetService.findAssetsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); |
|||
} |
|||
|
|||
@Override |
|||
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, Asset asset) { |
|||
return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ASSET, |
|||
EdgeEventActionType.ADDED, asset.getId(), null); |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue