433 changed files with 14199 additions and 1950 deletions
@ -0,0 +1,74 @@ |
|||
-- |
|||
-- Copyright © 2016-2022 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. |
|||
-- |
|||
|
|||
DO |
|||
$$ |
|||
DECLARE table_partition RECORD; |
|||
BEGIN |
|||
-- in case of running the upgrade script a second time: |
|||
IF NOT (SELECT exists(SELECT FROM pg_tables WHERE tablename = 'old_audit_log')) THEN |
|||
ALTER TABLE audit_log RENAME TO old_audit_log; |
|||
ALTER INDEX IF EXISTS idx_audit_log_tenant_id_and_created_time RENAME TO idx_old_audit_log_tenant_id_and_created_time; |
|||
|
|||
FOR table_partition IN SELECT tablename AS name, split_part(tablename, '_', 3) AS partition_ts |
|||
FROM pg_tables WHERE tablename LIKE 'audit_log_%' |
|||
LOOP |
|||
EXECUTE format('ALTER TABLE %s RENAME TO old_audit_log_%s', table_partition.name, table_partition.partition_ts); |
|||
END LOOP; |
|||
ELSE |
|||
RAISE NOTICE 'Table old_audit_log already exists, leaving as is'; |
|||
END IF; |
|||
END; |
|||
$$; |
|||
|
|||
CREATE TABLE IF NOT EXISTS audit_log ( |
|||
id uuid NOT NULL, |
|||
created_time bigint NOT NULL, |
|||
tenant_id uuid, |
|||
customer_id uuid, |
|||
entity_id uuid, |
|||
entity_type varchar(255), |
|||
entity_name varchar(255), |
|||
user_id uuid, |
|||
user_name varchar(255), |
|||
action_type varchar(255), |
|||
action_data varchar(1000000), |
|||
action_status varchar(255), |
|||
action_failure_details varchar(1000000) |
|||
) PARTITION BY RANGE (created_time); |
|||
CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_id_and_created_time ON audit_log(tenant_id, created_time DESC); |
|||
|
|||
CREATE OR REPLACE PROCEDURE migrate_audit_logs(IN start_time_ms BIGINT, IN end_time_ms BIGINT, IN partition_size_ms BIGINT) |
|||
LANGUAGE plpgsql AS |
|||
$$ |
|||
DECLARE |
|||
p RECORD; |
|||
partition_end_ts BIGINT; |
|||
BEGIN |
|||
FOR p IN SELECT DISTINCT (created_time - created_time % partition_size_ms) AS partition_ts FROM old_audit_log |
|||
WHERE created_time >= start_time_ms AND created_time < end_time_ms |
|||
LOOP |
|||
partition_end_ts = p.partition_ts + partition_size_ms; |
|||
RAISE NOTICE '[audit_log] Partition to create : [%-%]', p.partition_ts, partition_end_ts; |
|||
EXECUTE format('CREATE TABLE IF NOT EXISTS audit_log_%s PARTITION OF audit_log ' || |
|||
'FOR VALUES FROM ( %s ) TO ( %s )', p.partition_ts, p.partition_ts, partition_end_ts); |
|||
END LOOP; |
|||
|
|||
INSERT INTO audit_log |
|||
SELECT * FROM old_audit_log |
|||
WHERE created_time >= start_time_ms AND created_time < end_time_ms; |
|||
END; |
|||
$$; |
|||
@ -0,0 +1,21 @@ |
|||
-- |
|||
-- Copyright © 2016-2022 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. |
|||
-- |
|||
|
|||
DROP PROCEDURE IF EXISTS update_asset_profiles; |
|||
|
|||
ALTER TABLE asset ALTER COLUMN asset_profile_id SET NOT NULL; |
|||
ALTER TABLE asset DROP CONSTRAINT IF EXISTS fk_asset_profile; |
|||
ALTER TABLE asset ADD CONSTRAINT fk_asset_profile FOREIGN KEY (asset_profile_id) REFERENCES asset_profile(id); |
|||
@ -0,0 +1,45 @@ |
|||
-- |
|||
-- Copyright © 2016-2022 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. |
|||
-- |
|||
|
|||
CREATE TABLE IF NOT EXISTS asset_profile ( |
|||
id uuid NOT NULL CONSTRAINT asset_profile_pkey PRIMARY KEY, |
|||
created_time bigint NOT NULL, |
|||
name varchar(255), |
|||
image varchar(1000000), |
|||
description varchar, |
|||
search_text varchar(255), |
|||
is_default boolean, |
|||
tenant_id uuid, |
|||
default_rule_chain_id uuid, |
|||
default_dashboard_id uuid, |
|||
default_queue_name varchar(255), |
|||
external_id uuid, |
|||
CONSTRAINT asset_profile_name_unq_key UNIQUE (tenant_id, name), |
|||
CONSTRAINT asset_profile_external_id_unq_key UNIQUE (tenant_id, external_id), |
|||
CONSTRAINT fk_default_rule_chain_asset_profile FOREIGN KEY (default_rule_chain_id) REFERENCES rule_chain(id), |
|||
CONSTRAINT fk_default_dashboard_asset_profile FOREIGN KEY (default_dashboard_id) REFERENCES dashboard(id) |
|||
); |
|||
|
|||
CREATE OR REPLACE PROCEDURE update_asset_profiles() |
|||
LANGUAGE plpgsql AS |
|||
$$ |
|||
BEGIN |
|||
UPDATE asset as a SET asset_profile_id = p.id |
|||
FROM |
|||
(SELECT id, tenant_id, name from asset_profile) as p |
|||
WHERE a.asset_profile_id IS NULL AND p.tenant_id = a.tenant_id AND a.type = p.name; |
|||
END; |
|||
$$; |
|||
@ -0,0 +1,227 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.http.HttpStatus; |
|||
import org.springframework.security.access.prepost.PreAuthorize; |
|||
import org.springframework.web.bind.annotation.PathVariable; |
|||
import org.springframework.web.bind.annotation.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.ResponseStatus; |
|||
import org.springframework.web.bind.annotation.RestController; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.asset.AssetProfileInfo; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.entitiy.asset.profile.TbAssetProfileService; |
|||
import org.thingsboard.server.service.security.permission.Operation; |
|||
import org.thingsboard.server.service.security.permission.Resource; |
|||
|
|||
import static org.thingsboard.server.controller.ControllerConstants.ASSET_PROFILE_ID; |
|||
import static org.thingsboard.server.controller.ControllerConstants.ASSET_PROFILE_ID_PARAM_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.ASSET_PROFILE_INFO_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.ASSET_PROFILE_SORT_PROPERTY_ALLOWABLE_VALUES; |
|||
import static org.thingsboard.server.controller.ControllerConstants.ASSET_PROFILE_TEXT_SEARCH_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.NEW_LINE; |
|||
import static org.thingsboard.server.controller.ControllerConstants.PAGE_DATA_PARAMETERS; |
|||
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.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; |
|||
import static org.thingsboard.server.controller.ControllerConstants.UUID_WIKI_LINK; |
|||
|
|||
@RestController |
|||
@TbCoreComponent |
|||
@RequestMapping("/api") |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class AssetProfileController extends BaseController { |
|||
|
|||
private final TbAssetProfileService tbAssetProfileService; |
|||
|
|||
@ApiOperation(value = "Get Asset Profile (getAssetProfileById)", |
|||
notes = "Fetch the Asset Profile object based on the provided Asset Profile Id. " + |
|||
"The server checks that the asset profile is owned by the same tenant. " + TENANT_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/assetProfile/{assetProfileId}", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public AssetProfile getAssetProfileById( |
|||
@ApiParam(value = ASSET_PROFILE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable(ASSET_PROFILE_ID) String strAssetProfileId) throws ThingsboardException { |
|||
checkParameter(ASSET_PROFILE_ID, strAssetProfileId); |
|||
try { |
|||
AssetProfileId assetProfileId = new AssetProfileId(toUUID(strAssetProfileId)); |
|||
return checkAssetProfileId(assetProfileId, Operation.READ); |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
|
|||
@ApiOperation(value = "Get Asset Profile Info (getAssetProfileInfoById)", |
|||
notes = "Fetch the Asset Profile Info object based on the provided Asset Profile Id. " |
|||
+ ASSET_PROFILE_INFO_DESCRIPTION + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/assetProfileInfo/{assetProfileId}", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public AssetProfileInfo getAssetProfileInfoById( |
|||
@ApiParam(value = ASSET_PROFILE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable(ASSET_PROFILE_ID) String strAssetProfileId) throws ThingsboardException { |
|||
checkParameter(ASSET_PROFILE_ID, strAssetProfileId); |
|||
try { |
|||
AssetProfileId assetProfileId = new AssetProfileId(toUUID(strAssetProfileId)); |
|||
return new AssetProfileInfo(checkAssetProfileId(assetProfileId, Operation.READ)); |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
|
|||
@ApiOperation(value = "Get Default Asset Profile (getDefaultAssetProfileInfo)", |
|||
notes = "Fetch the Default Asset Profile Info object. " + |
|||
ASSET_PROFILE_INFO_DESCRIPTION + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/assetProfileInfo/default", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public AssetProfileInfo getDefaultAssetProfileInfo() throws ThingsboardException { |
|||
try { |
|||
return checkNotNull(assetProfileService.findDefaultAssetProfileInfo(getTenantId())); |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
|
|||
@ApiOperation(value = "Create Or Update Asset Profile (saveAssetProfile)", |
|||
notes = "Create or update the Asset Profile. When creating asset profile, platform generates asset profile id as " + UUID_WIKI_LINK + |
|||
"The newly created asset profile id will be present in the response. " + |
|||
"Specify existing asset profile id to update the asset profile. " + |
|||
"Referencing non-existing asset profile Id will cause 'Not Found' error. " + NEW_LINE + |
|||
"Asset profile name is unique in the scope of tenant. Only one 'default' asset profile may exist in scope of tenant. " + |
|||
"Remove 'id', 'tenantId' from the request body example (below) to create new Asset Profile entity. " + |
|||
TENANT_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json", |
|||
consumes = "application/json") |
|||
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/assetProfile", method = RequestMethod.POST) |
|||
@ResponseBody |
|||
public AssetProfile saveAssetProfile( |
|||
@ApiParam(value = "A JSON value representing the asset profile.") |
|||
@RequestBody AssetProfile assetProfile) throws Exception { |
|||
assetProfile.setTenantId(getTenantId()); |
|||
checkEntity(assetProfile.getId(), assetProfile, Resource.ASSET_PROFILE); |
|||
return tbAssetProfileService.save(assetProfile, getCurrentUser()); |
|||
} |
|||
|
|||
@ApiOperation(value = "Delete asset profile (deleteAssetProfile)", |
|||
notes = "Deletes the asset profile. Referencing non-existing asset profile Id will cause an error. " + |
|||
"Can't delete the asset profile if it is referenced by existing assets." + TENANT_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/assetProfile/{assetProfileId}", method = RequestMethod.DELETE) |
|||
@ResponseStatus(value = HttpStatus.OK) |
|||
public void deleteAssetProfile( |
|||
@ApiParam(value = ASSET_PROFILE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable(ASSET_PROFILE_ID) String strAssetProfileId) throws ThingsboardException { |
|||
checkParameter(ASSET_PROFILE_ID, strAssetProfileId); |
|||
AssetProfileId assetProfileId = new AssetProfileId(toUUID(strAssetProfileId)); |
|||
AssetProfile assetProfile = checkAssetProfileId(assetProfileId, Operation.DELETE); |
|||
tbAssetProfileService.delete(assetProfile, getCurrentUser()); |
|||
} |
|||
|
|||
@ApiOperation(value = "Make Asset Profile Default (setDefaultAssetProfile)", |
|||
notes = "Marks asset profile as default within a tenant scope." + TENANT_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/assetProfile/{assetProfileId}/default", method = RequestMethod.POST) |
|||
@ResponseBody |
|||
public AssetProfile setDefaultAssetProfile( |
|||
@ApiParam(value = ASSET_PROFILE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable(ASSET_PROFILE_ID) String strAssetProfileId) throws ThingsboardException { |
|||
checkParameter(ASSET_PROFILE_ID, strAssetProfileId); |
|||
AssetProfileId assetProfileId = new AssetProfileId(toUUID(strAssetProfileId)); |
|||
AssetProfile assetProfile = checkAssetProfileId(assetProfileId, Operation.WRITE); |
|||
AssetProfile previousDefaultAssetProfile = assetProfileService.findDefaultAssetProfile(getTenantId()); |
|||
return tbAssetProfileService.setDefaultAssetProfile(assetProfile, previousDefaultAssetProfile, getCurrentUser()); |
|||
} |
|||
|
|||
@ApiOperation(value = "Get Asset Profiles (getAssetProfiles)", |
|||
notes = "Returns a page of asset profile objects owned by tenant. " + |
|||
PAGE_DATA_PARAMETERS + TENANT_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/assetProfiles", params = {"pageSize", "page"}, method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public PageData<AssetProfile> getAssetProfiles( |
|||
@ApiParam(value = PAGE_SIZE_DESCRIPTION, required = true) |
|||
@RequestParam int pageSize, |
|||
@ApiParam(value = PAGE_NUMBER_DESCRIPTION, required = true) |
|||
@RequestParam int page, |
|||
@ApiParam(value = ASSET_PROFILE_TEXT_SEARCH_DESCRIPTION) |
|||
@RequestParam(required = false) String textSearch, |
|||
@ApiParam(value = SORT_PROPERTY_DESCRIPTION, allowableValues = ASSET_PROFILE_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 { |
|||
try { |
|||
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); |
|||
return checkNotNull(assetProfileService.findAssetProfiles(getTenantId(), pageLink)); |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
|
|||
@ApiOperation(value = "Get Asset Profile infos (getAssetProfileInfos)", |
|||
notes = "Returns a page of asset profile info objects owned by tenant. " + |
|||
PAGE_DATA_PARAMETERS + ASSET_PROFILE_INFO_DESCRIPTION + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH, |
|||
produces = "application/json") |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/assetProfileInfos", params = {"pageSize", "page"}, method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public PageData<AssetProfileInfo> getAssetProfileInfos( |
|||
@ApiParam(value = PAGE_SIZE_DESCRIPTION, required = true) |
|||
@RequestParam int pageSize, |
|||
@ApiParam(value = PAGE_NUMBER_DESCRIPTION, required = true) |
|||
@RequestParam int page, |
|||
@ApiParam(value = ASSET_PROFILE_TEXT_SEARCH_DESCRIPTION) |
|||
@RequestParam(required = false) String textSearch, |
|||
@ApiParam(value = SORT_PROPERTY_DESCRIPTION, allowableValues = ASSET_PROFILE_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 { |
|||
try { |
|||
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); |
|||
return checkNotNull(assetProfileService.findAssetProfileInfos(getTenantId(), pageLink)); |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,67 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.constructor; |
|||
|
|||
import com.google.protobuf.ByteString; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
import java.nio.charset.StandardCharsets; |
|||
|
|||
@Component |
|||
@TbCoreComponent |
|||
public class AssetProfileMsgConstructor { |
|||
|
|||
@Autowired |
|||
private DataDecodingEncodingService dataDecodingEncodingService; |
|||
|
|||
public AssetProfileUpdateMsg constructAssetProfileUpdatedMsg(UpdateMsgType msgType, AssetProfile assetProfile) { |
|||
AssetProfileUpdateMsg.Builder builder = AssetProfileUpdateMsg.newBuilder() |
|||
.setMsgType(msgType) |
|||
.setIdMSB(assetProfile.getId().getId().getMostSignificantBits()) |
|||
.setIdLSB(assetProfile.getId().getId().getLeastSignificantBits()) |
|||
.setName(assetProfile.getName()) |
|||
.setDefault(assetProfile.isDefault()); |
|||
if (assetProfile.getDefaultRuleChainId() != null) { |
|||
builder.setDefaultRuleChainIdMSB(assetProfile.getDefaultRuleChainId().getId().getMostSignificantBits()) |
|||
.setDefaultRuleChainIdLSB(assetProfile.getDefaultRuleChainId().getId().getLeastSignificantBits()); |
|||
} |
|||
if (assetProfile.getDefaultQueueName() != null) { |
|||
builder.setDefaultQueueName(assetProfile.getDefaultQueueName()); |
|||
} |
|||
if (assetProfile.getDescription() != null) { |
|||
builder.setDescription(assetProfile.getDescription()); |
|||
} |
|||
if (assetProfile.getImage() != null) { |
|||
builder.setImage(ByteString.copyFrom(assetProfile.getImage().getBytes(StandardCharsets.UTF_8))); |
|||
} |
|||
return builder.build(); |
|||
} |
|||
|
|||
public AssetProfileUpdateMsg constructAssetProfileDeleteMsg(AssetProfileId assetProfileId) { |
|||
return AssetProfileUpdateMsg.newBuilder() |
|||
.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) |
|||
.setIdMSB(assetProfileId.getId().getMostSignificantBits()) |
|||
.setIdLSB(assetProfileId.getId().getLeastSignificantBits()).build(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,47 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.EdgeUtils; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
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.AssetProfileService; |
|||
|
|||
@AllArgsConstructor |
|||
@Slf4j |
|||
public class AssetProfilesEdgeEventFetcher extends BasePageableEdgeEventFetcher<AssetProfile> { |
|||
|
|||
private final AssetProfileService assetProfileService; |
|||
|
|||
@Override |
|||
PageData<AssetProfile> fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { |
|||
return assetProfileService.findAssetProfiles(tenantId, pageLink); |
|||
} |
|||
|
|||
@Override |
|||
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, AssetProfile assetProfile) { |
|||
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ASSET_PROFILE, |
|||
EdgeEventActionType.ADDED, assetProfile.getId(), null); |
|||
} |
|||
} |
|||
@ -0,0 +1,63 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.processor; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.EdgeUtils; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
@Component |
|||
@Slf4j |
|||
@TbCoreComponent |
|||
public class AssetProfileEdgeProcessor extends BaseEdgeProcessor { |
|||
|
|||
public DownlinkMsg processAssetProfileToEdge(EdgeEvent edgeEvent, UpdateMsgType msgType, EdgeEventActionType action) { |
|||
AssetProfileId assetProfileId = new AssetProfileId(edgeEvent.getEntityId()); |
|||
DownlinkMsg downlinkMsg = null; |
|||
switch (action) { |
|||
case ADDED: |
|||
case UPDATED: |
|||
AssetProfile assetProfile = assetProfileService.findAssetProfileById(edgeEvent.getTenantId(), assetProfileId); |
|||
if (assetProfile != null) { |
|||
AssetProfileUpdateMsg assetProfileUpdateMsg = |
|||
assetProfileMsgConstructor.constructAssetProfileUpdatedMsg(msgType, assetProfile); |
|||
downlinkMsg = DownlinkMsg.newBuilder() |
|||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|||
.addAssetProfileUpdateMsg(assetProfileUpdateMsg) |
|||
.build(); |
|||
} |
|||
break; |
|||
case DELETED: |
|||
AssetProfileUpdateMsg assetProfileUpdateMsg = |
|||
assetProfileMsgConstructor.constructAssetProfileDeleteMsg(assetProfileId); |
|||
downlinkMsg = DownlinkMsg.newBuilder() |
|||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|||
.addAssetProfileUpdateMsg(assetProfileUpdateMsg) |
|||
.build(); |
|||
break; |
|||
} |
|||
return downlinkMsg; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,110 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.entitiy.asset.profile; |
|||
|
|||
import lombok.AllArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
|||
import org.thingsboard.server.dao.asset.AssetProfileService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.entitiy.AbstractTbEntityService; |
|||
|
|||
import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE; |
|||
|
|||
@Service |
|||
@TbCoreComponent |
|||
@AllArgsConstructor |
|||
@Slf4j |
|||
public class DefaultTbAssetProfileService extends AbstractTbEntityService implements TbAssetProfileService { |
|||
|
|||
private final AssetProfileService assetProfileService; |
|||
|
|||
@Override |
|||
public AssetProfile save(AssetProfile assetProfile, User user) throws Exception { |
|||
ActionType actionType = assetProfile.getId() == null ? ActionType.ADDED : ActionType.UPDATED; |
|||
TenantId tenantId = assetProfile.getTenantId(); |
|||
try { |
|||
if (TB_SERVICE_QUEUE.equals(assetProfile.getName())) { |
|||
throw new ThingsboardException("Unable to save asset profile with name " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS); |
|||
} else if (assetProfile.getId() != null) { |
|||
AssetProfile foundAssetProfile = assetProfileService.findAssetProfileById(tenantId, assetProfile.getId()); |
|||
if (foundAssetProfile != null && TB_SERVICE_QUEUE.equals(foundAssetProfile.getName())) { |
|||
throw new ThingsboardException("Updating asset profile with name " + TB_SERVICE_QUEUE + " is prohibited!", ThingsboardErrorCode.BAD_REQUEST_PARAMS); |
|||
} |
|||
} |
|||
AssetProfile savedAssetProfile = checkNotNull(assetProfileService.saveAssetProfile(assetProfile)); |
|||
autoCommit(user, savedAssetProfile.getId()); |
|||
tbClusterService.broadcastEntityStateChangeEvent(tenantId, savedAssetProfile.getId(), |
|||
actionType.equals(ActionType.ADDED) ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); |
|||
|
|||
notificationEntityService.notifyCreateOrUpdateOrDelete(tenantId, null, savedAssetProfile.getId(), |
|||
savedAssetProfile, user, actionType, true, null); |
|||
return savedAssetProfile; |
|||
} catch (Exception e) { |
|||
notificationEntityService.logEntityAction(tenantId, emptyId(EntityType.ASSET_PROFILE), assetProfile, actionType, user, e); |
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void delete(AssetProfile assetProfile, User user) { |
|||
AssetProfileId assetProfileId = assetProfile.getId(); |
|||
TenantId tenantId = assetProfile.getTenantId(); |
|||
try { |
|||
assetProfileService.deleteAssetProfile(tenantId, assetProfileId); |
|||
|
|||
tbClusterService.broadcastEntityStateChangeEvent(tenantId, assetProfileId, ComponentLifecycleEvent.DELETED); |
|||
notificationEntityService.notifyCreateOrUpdateOrDelete(tenantId, null, assetProfileId, assetProfile, |
|||
user, ActionType.DELETED, true, null, assetProfileId.toString()); |
|||
} catch (Exception e) { |
|||
notificationEntityService.logEntityAction(tenantId, emptyId(EntityType.ASSET_PROFILE), ActionType.DELETED, |
|||
user, e, assetProfileId.toString()); |
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public AssetProfile setDefaultAssetProfile(AssetProfile assetProfile, AssetProfile previousDefaultAssetProfile, User user) throws ThingsboardException { |
|||
TenantId tenantId = assetProfile.getTenantId(); |
|||
AssetProfileId assetProfileId = assetProfile.getId(); |
|||
try { |
|||
if (assetProfileService.setDefaultAssetProfile(tenantId, assetProfileId)) { |
|||
if (previousDefaultAssetProfile != null) { |
|||
previousDefaultAssetProfile = assetProfileService.findAssetProfileById(tenantId, previousDefaultAssetProfile.getId()); |
|||
notificationEntityService.logEntityAction(tenantId, previousDefaultAssetProfile.getId(), previousDefaultAssetProfile, |
|||
ActionType.UPDATED, user); |
|||
} |
|||
assetProfile = assetProfileService.findAssetProfileById(tenantId, assetProfileId); |
|||
|
|||
notificationEntityService.logEntityAction(tenantId, assetProfileId, assetProfile, ActionType.UPDATED, user); |
|||
} |
|||
return assetProfile; |
|||
} catch (Exception e) { |
|||
notificationEntityService.logEntityAction(tenantId, emptyId(EntityType.ASSET_PROFILE), ActionType.UPDATED, |
|||
user, e, assetProfileId.toString()); |
|||
throw e; |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,26 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.entitiy.asset.profile; |
|||
|
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.service.entitiy.SimpleTbEntityService; |
|||
|
|||
public interface TbAssetProfileService extends SimpleTbEntityService<AssetProfile> { |
|||
|
|||
AssetProfile setDefaultAssetProfile(AssetProfile assetProfile, AssetProfile previousDefaultAssetProfile, User user) throws ThingsboardException; |
|||
} |
|||
@ -0,0 +1,162 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.profile; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.dao.asset.AssetProfileService; |
|||
import org.thingsboard.server.dao.asset.AssetService; |
|||
|
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.ConcurrentMap; |
|||
import java.util.concurrent.locks.Lock; |
|||
import java.util.concurrent.locks.ReentrantLock; |
|||
import java.util.function.BiConsumer; |
|||
import java.util.function.Consumer; |
|||
|
|||
@Service |
|||
@Slf4j |
|||
public class DefaultTbAssetProfileCache implements TbAssetProfileCache { |
|||
|
|||
private final Lock assetProfileFetchLock = new ReentrantLock(); |
|||
private final AssetProfileService assetProfileService; |
|||
private final AssetService assetService; |
|||
|
|||
private final ConcurrentMap<AssetProfileId, AssetProfile> assetProfilesMap = new ConcurrentHashMap<>(); |
|||
private final ConcurrentMap<AssetId, AssetProfileId> assetsMap = new ConcurrentHashMap<>(); |
|||
private final ConcurrentMap<TenantId, ConcurrentMap<EntityId, Consumer<AssetProfile>>> profileListeners = new ConcurrentHashMap<>(); |
|||
private final ConcurrentMap<TenantId, ConcurrentMap<EntityId, BiConsumer<AssetId, AssetProfile>>> assetProfileListeners = new ConcurrentHashMap<>(); |
|||
|
|||
public DefaultTbAssetProfileCache(AssetProfileService assetProfileService, AssetService assetService) { |
|||
this.assetProfileService = assetProfileService; |
|||
this.assetService = assetService; |
|||
} |
|||
|
|||
@Override |
|||
public AssetProfile get(TenantId tenantId, AssetProfileId assetProfileId) { |
|||
AssetProfile profile = assetProfilesMap.get(assetProfileId); |
|||
if (profile == null) { |
|||
assetProfileFetchLock.lock(); |
|||
try { |
|||
profile = assetProfilesMap.get(assetProfileId); |
|||
if (profile == null) { |
|||
profile = assetProfileService.findAssetProfileById(tenantId, assetProfileId); |
|||
if (profile != null) { |
|||
assetProfilesMap.put(assetProfileId, profile); |
|||
log.debug("[{}] Fetch asset profile into cache: {}", profile.getId(), profile); |
|||
} |
|||
} |
|||
} finally { |
|||
assetProfileFetchLock.unlock(); |
|||
} |
|||
} |
|||
log.trace("[{}] Found asset profile in cache: {}", assetProfileId, profile); |
|||
return profile; |
|||
} |
|||
|
|||
@Override |
|||
public AssetProfile get(TenantId tenantId, AssetId assetId) { |
|||
AssetProfileId profileId = assetsMap.get(assetId); |
|||
if (profileId == null) { |
|||
Asset asset = assetService.findAssetById(tenantId, assetId); |
|||
if (asset != null) { |
|||
profileId = asset.getAssetProfileId(); |
|||
assetsMap.put(assetId, profileId); |
|||
} else { |
|||
return null; |
|||
} |
|||
} |
|||
return get(tenantId, profileId); |
|||
} |
|||
|
|||
@Override |
|||
public void evict(TenantId tenantId, AssetProfileId profileId) { |
|||
AssetProfile oldProfile = assetProfilesMap.remove(profileId); |
|||
log.debug("[{}] evict asset profile from cache: {}", profileId, oldProfile); |
|||
AssetProfile newProfile = get(tenantId, profileId); |
|||
if (newProfile != null) { |
|||
notifyProfileListeners(newProfile); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void evict(TenantId tenantId, AssetId assetId) { |
|||
AssetProfileId old = assetsMap.remove(assetId); |
|||
if (old != null) { |
|||
AssetProfile newProfile = get(tenantId, assetId); |
|||
if (newProfile == null || !old.equals(newProfile.getId())) { |
|||
notifyAssetListeners(tenantId, assetId, newProfile); |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void addListener(TenantId tenantId, EntityId listenerId, |
|||
Consumer<AssetProfile> profileListener, |
|||
BiConsumer<AssetId, AssetProfile> assetListener) { |
|||
if (profileListener != null) { |
|||
profileListeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, profileListener); |
|||
} |
|||
if (assetListener != null) { |
|||
assetProfileListeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, assetListener); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public AssetProfile find(AssetProfileId assetProfileId) { |
|||
return assetProfileService.findAssetProfileById(TenantId.SYS_TENANT_ID, assetProfileId); |
|||
} |
|||
|
|||
@Override |
|||
public AssetProfile findOrCreateAssetProfile(TenantId tenantId, String profileName) { |
|||
return assetProfileService.findOrCreateAssetProfile(tenantId, profileName); |
|||
} |
|||
|
|||
@Override |
|||
public void removeListener(TenantId tenantId, EntityId listenerId) { |
|||
ConcurrentMap<EntityId, Consumer<AssetProfile>> tenantListeners = profileListeners.get(tenantId); |
|||
if (tenantListeners != null) { |
|||
tenantListeners.remove(listenerId); |
|||
} |
|||
ConcurrentMap<EntityId, BiConsumer<AssetId, AssetProfile>> assetListeners = assetProfileListeners.get(tenantId); |
|||
if (assetListeners != null) { |
|||
assetListeners.remove(listenerId); |
|||
} |
|||
} |
|||
|
|||
private void notifyProfileListeners(AssetProfile profile) { |
|||
ConcurrentMap<EntityId, Consumer<AssetProfile>> tenantListeners = profileListeners.get(profile.getTenantId()); |
|||
if (tenantListeners != null) { |
|||
tenantListeners.forEach((id, listener) -> listener.accept(profile)); |
|||
} |
|||
} |
|||
|
|||
private void notifyAssetListeners(TenantId tenantId, AssetId assetId, AssetProfile profile) { |
|||
if (profile != null) { |
|||
ConcurrentMap<EntityId, BiConsumer<AssetId, AssetProfile>> tenantListeners = assetProfileListeners.get(tenantId); |
|||
if (tenantListeners != null) { |
|||
tenantListeners.forEach((id, listener) -> listener.accept(assetId, profile)); |
|||
} |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,33 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.profile; |
|||
|
|||
import org.thingsboard.rule.engine.api.RuleEngineAssetProfileCache; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
|
|||
public interface TbAssetProfileCache extends RuleEngineAssetProfileCache { |
|||
|
|||
void evict(TenantId tenantId, AssetProfileId id); |
|||
|
|||
void evict(TenantId tenantId, AssetId id); |
|||
|
|||
AssetProfile find(AssetProfileId assetProfileId); |
|||
|
|||
AssetProfile findOrCreateAssetProfile(TenantId tenantId, String assetType); |
|||
} |
|||
@ -0,0 +1,29 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.subscription; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
|||
import org.thingsboard.server.service.telemetry.cmd.v2.AggKey; |
|||
|
|||
@Data |
|||
public class ReadTsKvQueryInfo { |
|||
|
|||
private final AggKey key; |
|||
private final ReadTsKvQuery query; |
|||
private final boolean previous; |
|||
|
|||
} |
|||
@ -0,0 +1,43 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.sync.ie.exporting.impl; |
|||
|
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.sync.ie.EntityExportData; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.sync.vc.data.EntitiesExportCtx; |
|||
|
|||
import java.util.Set; |
|||
|
|||
@Service |
|||
@TbCoreComponent |
|||
public class AssetProfileExportService extends BaseEntityExportService<AssetProfileId, AssetProfile, EntityExportData<AssetProfile>> { |
|||
|
|||
@Override |
|||
protected void setRelatedEntities(EntitiesExportCtx<?> ctx, AssetProfile assetProfile, EntityExportData<AssetProfile> exportData) { |
|||
assetProfile.setDefaultDashboardId(getExternalIdOrElseInternal(ctx, assetProfile.getDefaultDashboardId())); |
|||
assetProfile.setDefaultRuleChainId(getExternalIdOrElseInternal(ctx, assetProfile.getDefaultRuleChainId())); |
|||
} |
|||
|
|||
@Override |
|||
public Set<EntityType> getSupportedEntityTypes() { |
|||
return Set.of(EntityType.ASSET_PROFILE); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,80 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.sync.ie.importing.impl; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
|||
import org.thingsboard.server.common.data.sync.ie.EntityExportData; |
|||
import org.thingsboard.server.dao.asset.AssetProfileService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.sync.vc.data.EntitiesImportCtx; |
|||
|
|||
@Service |
|||
@TbCoreComponent |
|||
@RequiredArgsConstructor |
|||
public class AssetProfileImportService extends BaseEntityImportService<AssetProfileId, AssetProfile, EntityExportData<AssetProfile>> { |
|||
|
|||
private final AssetProfileService assetProfileService; |
|||
|
|||
@Override |
|||
protected void setOwner(TenantId tenantId, AssetProfile assetProfile, IdProvider idProvider) { |
|||
assetProfile.setTenantId(tenantId); |
|||
} |
|||
|
|||
@Override |
|||
protected AssetProfile prepare(EntitiesImportCtx ctx, AssetProfile assetProfile, AssetProfile old, EntityExportData<AssetProfile> exportData, IdProvider idProvider) { |
|||
assetProfile.setDefaultRuleChainId(idProvider.getInternalId(assetProfile.getDefaultRuleChainId())); |
|||
assetProfile.setDefaultDashboardId(idProvider.getInternalId(assetProfile.getDefaultDashboardId())); |
|||
return assetProfile; |
|||
} |
|||
|
|||
@Override |
|||
protected AssetProfile saveOrUpdate(EntitiesImportCtx ctx, AssetProfile assetProfile, EntityExportData<AssetProfile> exportData, IdProvider idProvider) { |
|||
return assetProfileService.saveAssetProfile(assetProfile); |
|||
} |
|||
|
|||
@Override |
|||
protected void onEntitySaved(User user, AssetProfile savedAssetProfile, AssetProfile oldAssetProfile) throws ThingsboardException { |
|||
clusterService.broadcastEntityStateChangeEvent(user.getTenantId(), savedAssetProfile.getId(), |
|||
oldAssetProfile == null ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); |
|||
entityNotificationService.notifyCreateOrUpdateOrDelete(savedAssetProfile.getTenantId(), null, |
|||
savedAssetProfile.getId(), savedAssetProfile, user, oldAssetProfile == null ? ActionType.ADDED : ActionType.UPDATED, true, null); |
|||
} |
|||
|
|||
@Override |
|||
protected AssetProfile deepCopy(AssetProfile assetProfile) { |
|||
return new AssetProfile(assetProfile); |
|||
} |
|||
|
|||
@Override |
|||
protected void cleanupForComparison(AssetProfile assetProfile) { |
|||
super.cleanupForComparison(assetProfile); |
|||
} |
|||
|
|||
@Override |
|||
public EntityType getEntityType() { |
|||
return EntityType.ASSET_PROFILE; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,29 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.telemetry.cmd.v2; |
|||
|
|||
import lombok.Data; |
|||
|
|||
import java.util.List; |
|||
|
|||
@Data |
|||
public class AggHistoryCmd { |
|||
|
|||
private List<AggKey> keys; |
|||
private long startTs; |
|||
private long endTs; |
|||
|
|||
} |
|||
@ -0,0 +1,32 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.telemetry.cmd.v2; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.kv.Aggregation; |
|||
|
|||
@Data |
|||
public class AggKey { |
|||
|
|||
private int id; |
|||
private String key; |
|||
private Aggregation agg; |
|||
|
|||
private Long previousStartTs; |
|||
private Long previousEndTs; |
|||
private Boolean previousValueOnly; |
|||
|
|||
} |
|||
@ -0,0 +1,29 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.telemetry.cmd.v2; |
|||
|
|||
import lombok.Data; |
|||
|
|||
import java.util.List; |
|||
|
|||
@Data |
|||
public class AggTimeSeriesCmd { |
|||
|
|||
private List<AggKey> keys; |
|||
private long startTs; |
|||
private long timeWindow; |
|||
|
|||
} |
|||
@ -0,0 +1,61 @@ |
|||
/** |
|||
* Copyright © 2016-2022 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.ttl; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; |
|||
import org.springframework.scheduling.annotation.Scheduled; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.dao.audit.AuditLogDao; |
|||
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository; |
|||
import org.thingsboard.server.queue.discovery.PartitionService; |
|||
|
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
import static org.thingsboard.server.dao.model.ModelConstants.AUDIT_LOG_COLUMN_FAMILY_NAME; |
|||
|
|||
@Service |
|||
@ConditionalOnExpression("${sql.ttl.audit_logs.enabled:true} && ${sql.ttl.audit_logs.ttl:0} > 0") |
|||
@Slf4j |
|||
public class AuditLogsCleanUpService extends AbstractCleanUpService { |
|||
|
|||
private final AuditLogDao auditLogDao; |
|||
private final SqlPartitioningRepository partitioningRepository; |
|||
|
|||
@Value("${sql.ttl.audit_logs.ttl:0}") |
|||
private long ttlInSec; |
|||
@Value("${sql.audit_logs.partition_size:168}") |
|||
private int partitionSizeInHours; |
|||
|
|||
public AuditLogsCleanUpService(PartitionService partitionService, AuditLogDao auditLogDao, SqlPartitioningRepository partitioningRepository) { |
|||
super(partitionService); |
|||
this.auditLogDao = auditLogDao; |
|||
this.partitioningRepository = partitioningRepository; |
|||
} |
|||
|
|||
@Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.audit_logs.checking_interval_ms})}", |
|||
fixedDelayString = "${sql.ttl.audit_logs.checking_interval_ms}") |
|||
public void cleanUp() { |
|||
long auditLogsExpTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttlInSec); |
|||
if (isSystemTenantPartitionMine()) { |
|||
auditLogDao.cleanUpAuditLogs(auditLogsExpTime); |
|||
} else { |
|||
partitioningRepository.cleanupPartitionsCache(AUDIT_LOG_COLUMN_FAMILY_NAME, auditLogsExpTime, TimeUnit.HOURS.toMillis(partitionSizeInHours)); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,463 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.fasterxml.jackson.core.type.TypeReference; |
|||
import org.junit.After; |
|||
import org.junit.Assert; |
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.mockito.AdditionalAnswers; |
|||
import org.mockito.Mockito; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.context.annotation.Bean; |
|||
import org.springframework.context.annotation.Primary; |
|||
import org.springframework.test.context.ContextConfiguration; |
|||
import org.thingsboard.server.common.data.Customer; |
|||
import org.thingsboard.server.common.data.Dashboard; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.asset.AssetProfileInfo; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.common.data.rule.RuleChain; |
|||
import org.thingsboard.server.common.data.security.Authority; |
|||
import org.thingsboard.server.dao.asset.AssetProfileDao; |
|||
import org.thingsboard.server.dao.exception.DataValidationException; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.stream.Collectors; |
|||
|
|||
import static org.hamcrest.Matchers.containsString; |
|||
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; |
|||
|
|||
@ContextConfiguration(classes = {BaseAssetProfileControllerTest.Config.class}) |
|||
public abstract class BaseAssetProfileControllerTest extends AbstractControllerTest { |
|||
|
|||
private IdComparator<AssetProfile> idComparator = new IdComparator<>(); |
|||
private IdComparator<AssetProfileInfo> assetProfileInfoIdComparator = new IdComparator<>(); |
|||
|
|||
private Tenant savedTenant; |
|||
private User tenantAdmin; |
|||
|
|||
@Autowired |
|||
private AssetProfileDao assetProfileDao; |
|||
|
|||
static class Config { |
|||
@Bean |
|||
@Primary |
|||
public AssetProfileDao assetProfileDao(AssetProfileDao assetProfileDao) { |
|||
return Mockito.mock(AssetProfileDao.class, AdditionalAnswers.delegatesTo(assetProfileDao)); |
|||
} |
|||
} |
|||
|
|||
@Before |
|||
public void beforeTest() throws Exception { |
|||
loginSysAdmin(); |
|||
|
|||
Tenant tenant = new Tenant(); |
|||
tenant.setTitle("My tenant"); |
|||
savedTenant = doPost("/api/tenant", tenant, Tenant.class); |
|||
Assert.assertNotNull(savedTenant); |
|||
|
|||
tenantAdmin = new User(); |
|||
tenantAdmin.setAuthority(Authority.TENANT_ADMIN); |
|||
tenantAdmin.setTenantId(savedTenant.getId()); |
|||
tenantAdmin.setEmail("tenant2@thingsboard.org"); |
|||
tenantAdmin.setFirstName("Joe"); |
|||
tenantAdmin.setLastName("Downs"); |
|||
|
|||
tenantAdmin = createUserAndLogin(tenantAdmin, "testPassword1"); |
|||
} |
|||
|
|||
@After |
|||
public void afterTest() throws Exception { |
|||
loginSysAdmin(); |
|||
|
|||
doDelete("/api/tenant/" + savedTenant.getId().getId().toString()) |
|||
.andExpect(status().isOk()); |
|||
} |
|||
|
|||
@Test |
|||
public void testSaveAssetProfile() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
AssetProfile savedAssetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
Assert.assertNotNull(savedAssetProfile); |
|||
Assert.assertNotNull(savedAssetProfile.getId()); |
|||
Assert.assertTrue(savedAssetProfile.getCreatedTime() > 0); |
|||
Assert.assertEquals(assetProfile.getName(), savedAssetProfile.getName()); |
|||
Assert.assertEquals(assetProfile.getDescription(), savedAssetProfile.getDescription()); |
|||
Assert.assertEquals(assetProfile.isDefault(), savedAssetProfile.isDefault()); |
|||
Assert.assertEquals(assetProfile.getDefaultRuleChainId(), savedAssetProfile.getDefaultRuleChainId()); |
|||
|
|||
testNotifyEntityBroadcastEntityStateChangeEventOneTime(savedAssetProfile, savedAssetProfile.getId(), savedAssetProfile.getId(), |
|||
savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), |
|||
ActionType.ADDED); |
|||
|
|||
savedAssetProfile.setName("New asset profile"); |
|||
doPost("/api/assetProfile", savedAssetProfile, AssetProfile.class); |
|||
AssetProfile foundAssetProfile = doGet("/api/assetProfile/" + savedAssetProfile.getId().getId().toString(), AssetProfile.class); |
|||
Assert.assertEquals(savedAssetProfile.getName(), foundAssetProfile.getName()); |
|||
|
|||
testNotifyEntityBroadcastEntityStateChangeEventOneTime(foundAssetProfile, foundAssetProfile.getId(), foundAssetProfile.getId(), |
|||
savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), |
|||
ActionType.UPDATED); |
|||
} |
|||
|
|||
@Test |
|||
public void saveAssetProfileWithViolationOfValidation() throws Exception { |
|||
String msgError = msgErrorFieldLength("name"); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
AssetProfile createAssetProfile = this.createAssetProfile(StringUtils.randomAlphabetic(300)); |
|||
doPost("/api/assetProfile", createAssetProfile) |
|||
.andExpect(status().isBadRequest()) |
|||
.andExpect(statusReason(containsString(msgError))); |
|||
|
|||
testNotifyEntityEqualsOneTimeServiceNeverError(createAssetProfile, savedTenant.getId(), |
|||
tenantAdmin.getId(), tenantAdmin.getEmail(), ActionType.ADDED, new DataValidationException(msgError)); |
|||
} |
|||
|
|||
@Test |
|||
public void testFindAssetProfileById() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
AssetProfile savedAssetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
AssetProfile foundAssetProfile = doGet("/api/assetProfile/" + savedAssetProfile.getId().getId().toString(), AssetProfile.class); |
|||
Assert.assertNotNull(foundAssetProfile); |
|||
Assert.assertEquals(savedAssetProfile, foundAssetProfile); |
|||
} |
|||
|
|||
@Test |
|||
public void whenGetAssetProfileById_thenPermissionsAreChecked() throws Exception { |
|||
AssetProfile assetProfile = createAssetProfile("Asset profile 1"); |
|||
assetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
|
|||
loginDifferentTenant(); |
|||
|
|||
doGet("/api/assetProfile/" + assetProfile.getId()) |
|||
.andExpect(status().isForbidden()) |
|||
.andExpect(statusReason(containsString(msgErrorPermission))); |
|||
} |
|||
|
|||
@Test |
|||
public void testFindAssetProfileInfoById() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
AssetProfile savedAssetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
AssetProfileInfo foundAssetProfileInfo = doGet("/api/assetProfileInfo/" + savedAssetProfile.getId().getId().toString(), AssetProfileInfo.class); |
|||
Assert.assertNotNull(foundAssetProfileInfo); |
|||
Assert.assertEquals(savedAssetProfile.getId(), foundAssetProfileInfo.getId()); |
|||
Assert.assertEquals(savedAssetProfile.getName(), foundAssetProfileInfo.getName()); |
|||
|
|||
Customer customer = new Customer(); |
|||
customer.setTitle("Customer"); |
|||
customer.setTenantId(savedTenant.getId()); |
|||
Customer savedCustomer = doPost("/api/customer", customer, Customer.class); |
|||
|
|||
User customerUser = new User(); |
|||
customerUser.setAuthority(Authority.CUSTOMER_USER); |
|||
customerUser.setTenantId(savedTenant.getId()); |
|||
customerUser.setCustomerId(savedCustomer.getId()); |
|||
customerUser.setEmail("customer2@thingsboard.org"); |
|||
|
|||
createUserAndLogin(customerUser, "customer"); |
|||
|
|||
foundAssetProfileInfo = doGet("/api/assetProfileInfo/" + savedAssetProfile.getId().getId().toString(), AssetProfileInfo.class); |
|||
Assert.assertNotNull(foundAssetProfileInfo); |
|||
Assert.assertEquals(savedAssetProfile.getId(), foundAssetProfileInfo.getId()); |
|||
Assert.assertEquals(savedAssetProfile.getName(), foundAssetProfileInfo.getName()); |
|||
} |
|||
|
|||
@Test |
|||
public void whenGetAssetProfileInfoById_thenPermissionsAreChecked() throws Exception { |
|||
AssetProfile assetProfile = createAssetProfile("Asset profile 1"); |
|||
assetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
|
|||
loginDifferentTenant(); |
|||
doGet("/api/assetProfileInfo/" + assetProfile.getId()) |
|||
.andExpect(status().isForbidden()) |
|||
.andExpect(statusReason(containsString(msgErrorPermission))); |
|||
} |
|||
|
|||
@Test |
|||
public void testFindDefaultAssetProfileInfo() throws Exception { |
|||
AssetProfileInfo foundDefaultAssetProfileInfo = doGet("/api/assetProfileInfo/default", AssetProfileInfo.class); |
|||
Assert.assertNotNull(foundDefaultAssetProfileInfo); |
|||
Assert.assertNotNull(foundDefaultAssetProfileInfo.getId()); |
|||
Assert.assertNotNull(foundDefaultAssetProfileInfo.getName()); |
|||
Assert.assertEquals("default", foundDefaultAssetProfileInfo.getName()); |
|||
} |
|||
|
|||
@Test |
|||
public void testSetDefaultAssetProfile() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile 1"); |
|||
AssetProfile savedAssetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
AssetProfile defaultAssetProfile = doPost("/api/assetProfile/" + savedAssetProfile.getId().getId().toString() + "/default", AssetProfile.class); |
|||
Assert.assertNotNull(defaultAssetProfile); |
|||
AssetProfileInfo foundDefaultAssetProfile = doGet("/api/assetProfileInfo/default", AssetProfileInfo.class); |
|||
Assert.assertNotNull(foundDefaultAssetProfile); |
|||
Assert.assertEquals(savedAssetProfile.getName(), foundDefaultAssetProfile.getName()); |
|||
Assert.assertEquals(savedAssetProfile.getId(), foundDefaultAssetProfile.getId()); |
|||
|
|||
testNotifyEntityOneTimeMsgToEdgeServiceNever(defaultAssetProfile, defaultAssetProfile.getId(), defaultAssetProfile.getId(), |
|||
savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), |
|||
ActionType.UPDATED); |
|||
} |
|||
|
|||
@Test |
|||
public void testSaveAssetProfileWithEmptyName() throws Exception { |
|||
AssetProfile assetProfile = new AssetProfile(); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
String msgError = "Asset profile name " + msgErrorShouldBeSpecified; |
|||
doPost("/api/assetProfile", assetProfile) |
|||
.andExpect(status().isBadRequest()) |
|||
.andExpect(statusReason(containsString(msgError))); |
|||
|
|||
testNotifyEntityEqualsOneTimeServiceNeverError(assetProfile, savedTenant.getId(), |
|||
tenantAdmin.getId(), tenantAdmin.getEmail(), ActionType.ADDED, new DataValidationException(msgError)); |
|||
} |
|||
|
|||
@Test |
|||
public void testSaveAssetProfileWithSameName() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
doPost("/api/assetProfile", assetProfile).andExpect(status().isOk()); |
|||
AssetProfile assetProfile2 = this.createAssetProfile("Asset Profile"); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
String msgError = "Asset profile with such name already exists"; |
|||
doPost("/api/assetProfile", assetProfile2) |
|||
.andExpect(status().isBadRequest()) |
|||
.andExpect(statusReason(containsString(msgError))); |
|||
|
|||
testNotifyEntityEqualsOneTimeServiceNeverError(assetProfile, savedTenant.getId(), |
|||
tenantAdmin.getId(), tenantAdmin.getEmail(), ActionType.ADDED, new DataValidationException(msgError)); |
|||
} |
|||
|
|||
@Test |
|||
public void testDeleteAssetProfileWithExistingAsset() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
AssetProfile savedAssetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
|
|||
Asset asset = new Asset(); |
|||
asset.setName("Test asset"); |
|||
asset.setAssetProfileId(savedAssetProfile.getId()); |
|||
|
|||
doPost("/api/asset", asset, Asset.class); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
doDelete("/api/assetProfile/" + savedAssetProfile.getId().getId().toString()) |
|||
.andExpect(status().isBadRequest()) |
|||
.andExpect(statusReason(containsString("The asset profile referenced by the assets cannot be deleted"))); |
|||
|
|||
testNotifyEntityNever(savedAssetProfile.getId(), savedAssetProfile); |
|||
} |
|||
|
|||
@Test |
|||
public void testSaveAssetProfileWithRuleChainFromDifferentTenant() throws Exception { |
|||
loginDifferentTenant(); |
|||
RuleChain ruleChain = new RuleChain(); |
|||
ruleChain.setName("Different rule chain"); |
|||
RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class); |
|||
|
|||
loginTenantAdmin(); |
|||
|
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
assetProfile.setDefaultRuleChainId(savedRuleChain.getId()); |
|||
doPost("/api/assetProfile", assetProfile).andExpect(status().isBadRequest()) |
|||
.andExpect(statusReason(containsString("Can't assign rule chain from different tenant!"))); |
|||
} |
|||
|
|||
@Test |
|||
public void testSaveAssetProfileWithDashboardFromDifferentTenant() throws Exception { |
|||
loginDifferentTenant(); |
|||
Dashboard dashboard = new Dashboard(); |
|||
dashboard.setTitle("Different dashboard"); |
|||
Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class); |
|||
|
|||
loginTenantAdmin(); |
|||
|
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
assetProfile.setDefaultDashboardId(savedDashboard.getId()); |
|||
doPost("/api/assetProfile", assetProfile).andExpect(status().isBadRequest()) |
|||
.andExpect(statusReason(containsString("Can't assign dashboard from different tenant!"))); |
|||
} |
|||
|
|||
@Test |
|||
public void testDeleteAssetProfile() throws Exception { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile"); |
|||
AssetProfile savedAssetProfile = doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
doDelete("/api/assetProfile/" + savedAssetProfile.getId().getId().toString()) |
|||
.andExpect(status().isOk()); |
|||
|
|||
String savedAssetProfileIdStr = savedAssetProfile.getId().getId().toString(); |
|||
testNotifyEntityBroadcastEntityStateChangeEventOneTime(savedAssetProfile, savedAssetProfile.getId(), savedAssetProfile.getId(), |
|||
savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), |
|||
ActionType.DELETED, savedAssetProfileIdStr); |
|||
|
|||
doGet("/api/assetProfile/" + savedAssetProfile.getId().getId().toString()) |
|||
.andExpect(status().isNotFound()) |
|||
.andExpect(statusReason(containsString(msgErrorNoFound("Asset profile", savedAssetProfileIdStr)))); |
|||
} |
|||
|
|||
@Test |
|||
public void testFindAssetProfiles() throws Exception { |
|||
List<AssetProfile> assetProfiles = new ArrayList<>(); |
|||
PageLink pageLink = new PageLink(17); |
|||
PageData<AssetProfile> pageData = doGetTypedWithPageLink("/api/assetProfiles?", |
|||
new TypeReference<>() { |
|||
}, pageLink); |
|||
Assert.assertFalse(pageData.hasNext()); |
|||
Assert.assertEquals(1, pageData.getTotalElements()); |
|||
assetProfiles.addAll(pageData.getData()); |
|||
|
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
int cntEntity = 28; |
|||
for (int i = 0; i < cntEntity; i++) { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile" + i); |
|||
assetProfiles.add(doPost("/api/assetProfile", assetProfile, AssetProfile.class)); |
|||
} |
|||
|
|||
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(new AssetProfile(), new AssetProfile(), |
|||
savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), |
|||
ActionType.ADDED, ActionType.ADDED, cntEntity, cntEntity, cntEntity); |
|||
Mockito.reset(tbClusterService, auditLogService); |
|||
|
|||
List<AssetProfile> loadedAssetProfiles = new ArrayList<>(); |
|||
pageLink = new PageLink(17); |
|||
do { |
|||
pageData = doGetTypedWithPageLink("/api/assetProfiles?", |
|||
new TypeReference<>() { |
|||
}, pageLink); |
|||
loadedAssetProfiles.addAll(pageData.getData()); |
|||
if (pageData.hasNext()) { |
|||
pageLink = pageLink.nextPageLink(); |
|||
} |
|||
} while (pageData.hasNext()); |
|||
|
|||
Collections.sort(assetProfiles, idComparator); |
|||
Collections.sort(loadedAssetProfiles, idComparator); |
|||
|
|||
Assert.assertEquals(assetProfiles, loadedAssetProfiles); |
|||
|
|||
for (AssetProfile assetProfile : loadedAssetProfiles) { |
|||
if (!assetProfile.isDefault()) { |
|||
doDelete("/api/assetProfile/" + assetProfile.getId().getId().toString()) |
|||
.andExpect(status().isOk()); |
|||
} |
|||
} |
|||
|
|||
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(loadedAssetProfiles.get(0), loadedAssetProfiles.get(0), |
|||
savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), |
|||
ActionType.DELETED, ActionType.DELETED, cntEntity, cntEntity, cntEntity, loadedAssetProfiles.get(0).getId().getId().toString()); |
|||
|
|||
pageLink = new PageLink(17); |
|||
pageData = doGetTypedWithPageLink("/api/assetProfiles?", |
|||
new TypeReference<>() { |
|||
}, pageLink); |
|||
Assert.assertFalse(pageData.hasNext()); |
|||
Assert.assertEquals(1, pageData.getTotalElements()); |
|||
} |
|||
|
|||
@Test |
|||
public void testFindAssetProfileInfos() throws Exception { |
|||
List<AssetProfile> assetProfiles = new ArrayList<>(); |
|||
PageLink pageLink = new PageLink(17); |
|||
PageData<AssetProfile> assetProfilePageData = doGetTypedWithPageLink("/api/assetProfiles?", |
|||
new TypeReference<PageData<AssetProfile>>() { |
|||
}, pageLink); |
|||
Assert.assertFalse(assetProfilePageData.hasNext()); |
|||
Assert.assertEquals(1, assetProfilePageData.getTotalElements()); |
|||
assetProfiles.addAll(assetProfilePageData.getData()); |
|||
|
|||
for (int i = 0; i < 28; i++) { |
|||
AssetProfile assetProfile = this.createAssetProfile("Asset Profile" + i); |
|||
assetProfiles.add(doPost("/api/assetProfile", assetProfile, AssetProfile.class)); |
|||
} |
|||
|
|||
List<AssetProfileInfo> loadedAssetProfileInfos = new ArrayList<>(); |
|||
pageLink = new PageLink(17); |
|||
PageData<AssetProfileInfo> pageData; |
|||
do { |
|||
pageData = doGetTypedWithPageLink("/api/assetProfileInfos?", |
|||
new TypeReference<>() { |
|||
}, pageLink); |
|||
loadedAssetProfileInfos.addAll(pageData.getData()); |
|||
if (pageData.hasNext()) { |
|||
pageLink = pageLink.nextPageLink(); |
|||
} |
|||
} while (pageData.hasNext()); |
|||
|
|||
Collections.sort(assetProfiles, idComparator); |
|||
Collections.sort(loadedAssetProfileInfos, assetProfileInfoIdComparator); |
|||
|
|||
List<AssetProfileInfo> assetProfileInfos = assetProfiles.stream().map(assetProfile -> new AssetProfileInfo(assetProfile.getId(), |
|||
assetProfile.getName(), assetProfile.getImage(), assetProfile.getDefaultDashboardId())).collect(Collectors.toList()); |
|||
|
|||
Assert.assertEquals(assetProfileInfos, loadedAssetProfileInfos); |
|||
|
|||
for (AssetProfile assetProfile : assetProfiles) { |
|||
if (!assetProfile.isDefault()) { |
|||
doDelete("/api/assetProfile/" + assetProfile.getId().getId().toString()) |
|||
.andExpect(status().isOk()); |
|||
} |
|||
} |
|||
|
|||
pageLink = new PageLink(17); |
|||
pageData = doGetTypedWithPageLink("/api/assetProfileInfos?", |
|||
new TypeReference<PageData<AssetProfileInfo>>() { |
|||
}, pageLink); |
|||
Assert.assertFalse(pageData.hasNext()); |
|||
Assert.assertEquals(1, pageData.getTotalElements()); |
|||
} |
|||
|
|||
@Test |
|||
public void testDeleteAssetProfileWithDeleteRelationsOk() throws Exception { |
|||
AssetProfileId assetProfileId = savedAssetProfile("AssetProfile for Test WithRelationsOk").getId(); |
|||
testEntityDaoWithRelationsOk(savedTenant.getId(), assetProfileId, "/api/assetProfile/" + assetProfileId); |
|||
} |
|||
|
|||
@Test |
|||
public void testDeleteAssetProfileExceptionWithRelationsTransactional() throws Exception { |
|||
AssetProfileId assetProfileId = savedAssetProfile("AssetProfile for Test WithRelations Transactional Exception").getId(); |
|||
testEntityDaoWithRelationsTransactionalException(assetProfileDao, savedTenant.getId(), assetProfileId, "/api/assetProfile/" + assetProfileId); |
|||
} |
|||
|
|||
private AssetProfile savedAssetProfile(String name) { |
|||
AssetProfile assetProfile = createAssetProfile(name); |
|||
return doPost("/api/assetProfile", assetProfile, AssetProfile.class); |
|||
} |
|||
} |
|||
@ -0,0 +1,23 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.sql; |
|||
|
|||
import org.thingsboard.server.controller.BaseAssetProfileControllerTest; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
@DaoSqlTest |
|||
public class AssetProfileControllerSqlTest extends BaseAssetProfileControllerTest { |
|||
} |
|||
@ -0,0 +1,96 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.script; |
|||
|
|||
import org.junit.jupiter.api.Test; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.test.context.TestPropertySource; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.controller.AbstractControllerTest; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
import java.util.UUID; |
|||
import java.util.concurrent.ExecutionException; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
|||
|
|||
@DaoSqlTest |
|||
@TestPropertySource(properties = { |
|||
"js.max_script_body_size=50", |
|||
"js.max_total_args_size=50", |
|||
"js.max_result_size=50", |
|||
"js.local.max_errors=2" |
|||
}) |
|||
class LocalJsInvokeServiceTest extends AbstractControllerTest { |
|||
|
|||
@Autowired |
|||
private NashornJsInvokeService jsInvokeService; |
|||
|
|||
@Value("${js.local.max_errors}") |
|||
private int maxJsErrors; |
|||
|
|||
@Test |
|||
void givenTooBigScriptForEval_thenReturnError() { |
|||
String hugeScript = "var a = 'qwertyqwertywertyqwabababer'; return {a: a};"; |
|||
|
|||
assertThatThrownBy(() -> { |
|||
evalScript(hugeScript); |
|||
}).hasMessageContaining("body exceeds maximum allowed size"); |
|||
} |
|||
|
|||
@Test |
|||
void givenTooBigScriptInputArgs_thenReturnErrorAndReportScriptExecutionError() throws Exception { |
|||
String script = "return { msg: msg };"; |
|||
String hugeMsg = "{\"input\":\"123456781234349\"}"; |
|||
UUID scriptId = evalScript(script); |
|||
|
|||
for (int i = 0; i < maxJsErrors; i++) { |
|||
assertThatThrownBy(() -> { |
|||
invokeScript(scriptId, hugeMsg); |
|||
}).hasMessageContaining("input arguments exceed maximum"); |
|||
} |
|||
assertThatScriptIsBlocked(scriptId); |
|||
} |
|||
|
|||
@Test |
|||
void whenScriptInvocationResultIsTooBig_thenReturnErrorAndReportScriptExecutionError() throws Exception { |
|||
String script = "var s = new Array(50).join('a'); return { s: s};"; |
|||
UUID scriptId = evalScript(script); |
|||
|
|||
for (int i = 0; i < maxJsErrors; i++) { |
|||
assertThatThrownBy(() -> { |
|||
invokeScript(scriptId, "{}"); |
|||
}).hasMessageContaining("result exceeds maximum allowed size"); |
|||
} |
|||
assertThatScriptIsBlocked(scriptId); |
|||
} |
|||
|
|||
private void assertThatScriptIsBlocked(UUID scriptId) { |
|||
assertThatThrownBy(() -> { |
|||
invokeScript(scriptId, "{}"); |
|||
}).hasMessageContaining("invocation is blocked due to maximum error"); |
|||
} |
|||
|
|||
private UUID evalScript(String script) throws ExecutionException, InterruptedException { |
|||
return jsInvokeService.eval(TenantId.SYS_TENANT_ID, JsScriptType.RULE_NODE_SCRIPT, script).get(); |
|||
} |
|||
|
|||
private String invokeScript(UUID scriptId, String msg) throws ExecutionException, InterruptedException { |
|||
return jsInvokeService.invokeFunction(TenantId.SYS_TENANT_ID, null, scriptId, msg, "{}", "POST_TELEMETRY_REQUEST").get(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,219 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.script; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.junit.jupiter.api.AfterEach; |
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.thingsboard.server.common.data.ApiUsageState; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.js.JsInvokeProtos; |
|||
import org.thingsboard.server.gen.js.JsInvokeProtos.RemoteJsRequest; |
|||
import org.thingsboard.server.gen.js.JsInvokeProtos.RemoteJsResponse; |
|||
import org.thingsboard.server.queue.TbQueueRequestTemplate; |
|||
import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; |
|||
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
|||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|||
|
|||
import java.util.HashSet; |
|||
import java.util.List; |
|||
import java.util.Set; |
|||
import java.util.UUID; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.argThat; |
|||
import static org.mockito.Mockito.doAnswer; |
|||
import static org.mockito.Mockito.doReturn; |
|||
import static org.mockito.Mockito.mock; |
|||
import static org.mockito.Mockito.reset; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.verifyNoInteractions; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
class RemoteJsInvokeServiceTest { |
|||
|
|||
private RemoteJsInvokeService remoteJsInvokeService; |
|||
private TbQueueRequestTemplate<TbProtoJsQueueMsg<RemoteJsRequest>, TbProtoQueueMsg<RemoteJsResponse>> jsRequestTemplate; |
|||
|
|||
|
|||
@BeforeEach |
|||
public void beforeEach() { |
|||
TbApiUsageStateService apiUsageStateService = mock(TbApiUsageStateService.class); |
|||
ApiUsageState apiUsageState = mock(ApiUsageState.class); |
|||
when(apiUsageState.isJsExecEnabled()).thenReturn(true); |
|||
when(apiUsageStateService.getApiUsageState(any())).thenReturn(apiUsageState); |
|||
TbApiUsageClient apiUsageClient = mock(TbApiUsageClient.class); |
|||
|
|||
remoteJsInvokeService = new RemoteJsInvokeService(apiUsageStateService, apiUsageClient); |
|||
jsRequestTemplate = mock(TbQueueRequestTemplate.class); |
|||
remoteJsInvokeService.requestTemplate = jsRequestTemplate; |
|||
} |
|||
|
|||
@AfterEach |
|||
public void afterEach() { |
|||
reset(jsRequestTemplate); |
|||
} |
|||
|
|||
@Test |
|||
public void whenInvokingFunction_thenDoNotSendScriptBody() throws Exception { |
|||
mockJsEvalResponse(); |
|||
String scriptBody = "return { a: 'b'};"; |
|||
UUID scriptId = remoteJsInvokeService.eval(TenantId.SYS_TENANT_ID, JsScriptType.RULE_NODE_SCRIPT, scriptBody).get(); |
|||
reset(jsRequestTemplate); |
|||
|
|||
String expectedInvocationResult = "scriptInvocationResult"; |
|||
doReturn(Futures.immediateFuture(new TbProtoJsQueueMsg<>(UUID.randomUUID(), RemoteJsResponse.newBuilder() |
|||
.setInvokeResponse(JsInvokeProtos.JsInvokeResponse.newBuilder() |
|||
.setSuccess(true) |
|||
.setResult(expectedInvocationResult) |
|||
.build()) |
|||
.build()))) |
|||
.when(jsRequestTemplate).send(any()); |
|||
|
|||
ArgumentCaptor<TbProtoJsQueueMsg<RemoteJsRequest>> jsRequestCaptor = ArgumentCaptor.forClass(TbProtoJsQueueMsg.class); |
|||
Object invocationResult = remoteJsInvokeService.invokeFunction(TenantId.SYS_TENANT_ID, null, scriptId, "{}").get(); |
|||
verify(jsRequestTemplate).send(jsRequestCaptor.capture()); |
|||
|
|||
JsInvokeProtos.JsInvokeRequest jsInvokeRequestMade = jsRequestCaptor.getValue().getValue().getInvokeRequest(); |
|||
assertThat(jsInvokeRequestMade.getScriptBody()).isNullOrEmpty(); |
|||
assertThat(jsInvokeRequestMade.getScriptHash()).isEqualTo(getScriptHash(scriptId)); |
|||
assertThat(invocationResult).isEqualTo(expectedInvocationResult); |
|||
} |
|||
|
|||
@Test |
|||
public void whenInvokingFunctionAndRemoteJsExecutorRemovedScript_thenHandleNotFoundErrorAndMakeInvokeRequestWithScriptBody() throws Exception { |
|||
mockJsEvalResponse(); |
|||
String scriptBody = "return { a: 'b'};"; |
|||
UUID scriptId = remoteJsInvokeService.eval(TenantId.SYS_TENANT_ID, JsScriptType.RULE_NODE_SCRIPT, scriptBody).get(); |
|||
reset(jsRequestTemplate); |
|||
|
|||
doReturn(Futures.immediateFuture(new TbProtoJsQueueMsg<>(UUID.randomUUID(), RemoteJsResponse.newBuilder() |
|||
.setInvokeResponse(JsInvokeProtos.JsInvokeResponse.newBuilder() |
|||
.setSuccess(false) |
|||
.setErrorCode(JsInvokeProtos.JsInvokeErrorCode.NOT_FOUND_ERROR) |
|||
.build()) |
|||
.build()))) |
|||
.when(jsRequestTemplate).send(argThat(jsQueueMsg -> { |
|||
return StringUtils.isEmpty(jsQueueMsg.getValue().getInvokeRequest().getScriptBody()); |
|||
})); |
|||
|
|||
String expectedInvocationResult = "invocationResult"; |
|||
doReturn(Futures.immediateFuture(new TbProtoJsQueueMsg<>(UUID.randomUUID(), RemoteJsResponse.newBuilder() |
|||
.setInvokeResponse(JsInvokeProtos.JsInvokeResponse.newBuilder() |
|||
.setSuccess(true) |
|||
.setResult(expectedInvocationResult) |
|||
.build()) |
|||
.build()))) |
|||
.when(jsRequestTemplate).send(argThat(jsQueueMsg -> { |
|||
return StringUtils.isNotEmpty(jsQueueMsg.getValue().getInvokeRequest().getScriptBody()); |
|||
})); |
|||
|
|||
ArgumentCaptor<TbProtoJsQueueMsg<RemoteJsRequest>> jsRequestsCaptor = ArgumentCaptor.forClass(TbProtoJsQueueMsg.class); |
|||
Object invocationResult = remoteJsInvokeService.invokeFunction(TenantId.SYS_TENANT_ID, null, scriptId, "{}").get(); |
|||
verify(jsRequestTemplate, times(2)).send(jsRequestsCaptor.capture()); |
|||
|
|||
List<TbProtoJsQueueMsg<RemoteJsRequest>> jsInvokeRequestsMade = jsRequestsCaptor.getAllValues(); |
|||
|
|||
JsInvokeProtos.JsInvokeRequest firstRequestMade = jsInvokeRequestsMade.get(0).getValue().getInvokeRequest(); |
|||
assertThat(firstRequestMade.getScriptBody()).isNullOrEmpty(); |
|||
|
|||
JsInvokeProtos.JsInvokeRequest secondRequestMade = jsInvokeRequestsMade.get(1).getValue().getInvokeRequest(); |
|||
assertThat(secondRequestMade.getScriptBody()).contains(scriptBody); |
|||
|
|||
assertThat(jsInvokeRequestsMade.stream().map(TbProtoQueueMsg::getKey).distinct().count()).as("partition keys are same") |
|||
.isOne(); |
|||
|
|||
assertThat(invocationResult).isEqualTo(expectedInvocationResult); |
|||
} |
|||
|
|||
@Test |
|||
public void whenDoingEval_thenSaveScriptByHashOfTenantIdAndScriptBody() throws Exception { |
|||
mockJsEvalResponse(); |
|||
|
|||
TenantId tenantId1 = TenantId.fromUUID(UUID.randomUUID()); |
|||
String scriptBody1 = "var msg = { temp: 42, humidity: 77 };\n" + |
|||
"var metadata = { data: 40 };\n" + |
|||
"var msgType = \"POST_TELEMETRY_REQUEST\";\n" + |
|||
"\n" + |
|||
"return { msg: msg, metadata: metadata, msgType: msgType };"; |
|||
|
|||
Set<String> scriptHashes = new HashSet<>(); |
|||
String tenant1Script1Hash = null; |
|||
for (int i = 0; i < 3; i++) { |
|||
UUID scriptUuid = remoteJsInvokeService.eval(tenantId1, JsScriptType.RULE_NODE_SCRIPT, scriptBody1).get(); |
|||
tenant1Script1Hash = getScriptHash(scriptUuid); |
|||
scriptHashes.add(tenant1Script1Hash); |
|||
} |
|||
assertThat(scriptHashes).as("Unique scripts ids").size().isOne(); |
|||
|
|||
TenantId tenantId2 = TenantId.fromUUID(UUID.randomUUID()); |
|||
UUID scriptUuid = remoteJsInvokeService.eval(tenantId2, JsScriptType.RULE_NODE_SCRIPT, scriptBody1).get(); |
|||
String tenant2Script1Id = getScriptHash(scriptUuid); |
|||
assertThat(tenant2Script1Id).isNotEqualTo(tenant1Script1Hash); |
|||
|
|||
String scriptBody2 = scriptBody1 + ";;"; |
|||
scriptUuid = remoteJsInvokeService.eval(tenantId2, JsScriptType.RULE_NODE_SCRIPT, scriptBody2).get(); |
|||
String tenant2Script2Id = getScriptHash(scriptUuid); |
|||
assertThat(tenant2Script2Id).isNotEqualTo(tenant2Script1Id); |
|||
} |
|||
|
|||
@Test |
|||
public void whenReleasingScript_thenCheckForHashUsages() throws Exception { |
|||
mockJsEvalResponse(); |
|||
String scriptBody = "return { a: 'b'};"; |
|||
UUID scriptId1 = remoteJsInvokeService.eval(TenantId.SYS_TENANT_ID, JsScriptType.RULE_NODE_SCRIPT, scriptBody).get(); |
|||
UUID scriptId2 = remoteJsInvokeService.eval(TenantId.SYS_TENANT_ID, JsScriptType.RULE_NODE_SCRIPT, scriptBody).get(); |
|||
String scriptHash = getScriptHash(scriptId1); |
|||
assertThat(scriptHash).isEqualTo(getScriptHash(scriptId2)); |
|||
reset(jsRequestTemplate); |
|||
|
|||
doReturn(Futures.immediateFuture(new TbProtoQueueMsg<>(UUID.randomUUID(), RemoteJsResponse.newBuilder() |
|||
.setReleaseResponse(JsInvokeProtos.JsReleaseResponse.newBuilder() |
|||
.setSuccess(true) |
|||
.build()) |
|||
.build()))) |
|||
.when(jsRequestTemplate).send(any()); |
|||
|
|||
remoteJsInvokeService.release(scriptId1).get(); |
|||
verifyNoInteractions(jsRequestTemplate); |
|||
assertThat(remoteJsInvokeService.scriptHashToBodysMap).containsKey(scriptHash); |
|||
|
|||
remoteJsInvokeService.release(scriptId2).get(); |
|||
verify(jsRequestTemplate).send(any()); |
|||
assertThat(remoteJsInvokeService.scriptHashToBodysMap).isEmpty(); |
|||
} |
|||
|
|||
private String getScriptHash(UUID scriptUuid) { |
|||
return remoteJsInvokeService.scriptIdToNameAndHashMap.get(scriptUuid).getSecond(); |
|||
} |
|||
|
|||
private void mockJsEvalResponse() { |
|||
doAnswer(methodCall -> Futures.immediateFuture(new TbProtoJsQueueMsg<>(UUID.randomUUID(), RemoteJsResponse.newBuilder() |
|||
.setCompileResponse(JsInvokeProtos.JsCompileResponse.newBuilder() |
|||
.setSuccess(true) |
|||
.setScriptHash(methodCall.<TbProtoQueueMsg<RemoteJsRequest>>getArgument(0).getValue().getCompileRequest().getScriptHash()) |
|||
.build()) |
|||
.build()))) |
|||
.when(jsRequestTemplate).send(argThat(jsQueueMsg -> jsQueueMsg.getValue().hasCompileRequest())); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,53 @@ |
|||
/** |
|||
* Copyright © 2016-2022 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.dao.asset; |
|||
|
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.asset.AssetProfileInfo; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
|
|||
public interface AssetProfileService { |
|||
|
|||
AssetProfile findAssetProfileById(TenantId tenantId, AssetProfileId assetProfileId); |
|||
|
|||
AssetProfile findAssetProfileByName(TenantId tenantId, String profileName); |
|||
|
|||
AssetProfileInfo findAssetProfileInfoById(TenantId tenantId, AssetProfileId assetProfileId); |
|||
|
|||
AssetProfile saveAssetProfile(AssetProfile assetProfile); |
|||
|
|||
void deleteAssetProfile(TenantId tenantId, AssetProfileId assetProfileId); |
|||
|
|||
PageData<AssetProfile> findAssetProfiles(TenantId tenantId, PageLink pageLink); |
|||
|
|||
PageData<AssetProfileInfo> findAssetProfileInfos(TenantId tenantId, PageLink pageLink); |
|||
|
|||
AssetProfile findOrCreateAssetProfile(TenantId tenantId, String profileName); |
|||
|
|||
AssetProfile createDefaultAssetProfile(TenantId tenantId); |
|||
|
|||
AssetProfile findDefaultAssetProfile(TenantId tenantId); |
|||
|
|||
AssetProfileInfo findDefaultAssetProfileInfo(TenantId tenantId); |
|||
|
|||
boolean setDefaultAssetProfile(TenantId tenantId, AssetProfileId assetProfileId); |
|||
|
|||
void deleteAssetProfilesByTenantId(TenantId tenantId); |
|||
|
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue