From 43c2ae653d2b47d2725d54e6ba1cb633d1c0f5a4 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 11 Sep 2023 17:04:25 +0300 Subject: [PATCH 1/5] Resource support for edge --- .../edge/DefaultEdgeNotificationService.java | 7 ++ .../service/edge/EdgeContextComponent.java | 8 ++ .../service/edge/rpc/EdgeGrpcSession.java | 8 ++ .../service/edge/rpc/EdgeSyncCursor.java | 4 + .../constructor/ResourceMsgConstructor.java | 57 ++++++++++ .../fetch/BaseResourceEdgeEventFetcher.java | 49 +++++++++ .../SystemResourcesEdgeEventFetcher.java | 34 ++++++ .../TenantResourcesEdgeEventFetcher.java | 36 +++++++ .../TenantWidgetsBundlesEdgeEventFetcher.java | 1 + .../edge/rpc/processor/BaseEdgeProcessor.java | 12 +++ .../resource/BaseResourceProcessor.java | 61 +++++++++++ .../resource/ResourceEdgeProcessor.java | 100 ++++++++++++++++++ .../server/edge/AbstractEdgeTest.java | 3 +- .../asset/AssetEdgeProcessorTest.java | 3 +- .../asset/AssetProfileEdgeProcessorTest.java | 2 +- .../device/DeviceEdgeProcessorTest.java | 3 +- .../DeviceProfileEdgeProcessorTest.java | 3 +- .../server/dao/resource/ResourceService.java | 4 + .../common/data/edge/EdgeEventType.java | 3 +- .../common/data/id/EntityIdFactory.java | 2 + common/edge-api/src/main/proto/edge.proto | 15 +++ .../dao/resource/BaseResourceService.java | 28 ++++- 22 files changed, 430 insertions(+), 13 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/ResourceMsgConstructor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseResourceEdgeEventFetcher.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/SystemResourcesEdgeEventFetcher.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantResourcesEdgeEventFetcher.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/BaseResourceProcessor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index 69852e6c05..f3832d8fcf 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -48,6 +48,7 @@ import org.thingsboard.server.service.edge.rpc.processor.entityview.EntityViewEd import org.thingsboard.server.service.edge.rpc.processor.ota.OtaPackageEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.queue.QueueEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.relation.RelationEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.resource.ResourceEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.rule.RuleChainEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.tenant.TenantEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.tenant.TenantProfileEdgeProcessor; @@ -125,6 +126,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { @Autowired private RelationEdgeProcessor relationProcessor; + @Autowired + private ResourceEdgeProcessor resourceEdgeProcessor; + @Autowired protected ApplicationEventPublisher eventPublisher; @@ -215,6 +219,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { case TENANT_PROFILE: future = tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break; + case TB_RESOURCE: + future = resourceEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; default: log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); future = Futures.immediateFuture(null); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index dc83e0f2a2..7850032142 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -33,6 +33,7 @@ import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; +import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.settings.AdminSettingsService; import org.thingsboard.server.dao.tenant.TenantProfileService; @@ -55,6 +56,7 @@ import org.thingsboard.server.service.edge.rpc.processor.entityview.EntityViewEd import org.thingsboard.server.service.edge.rpc.processor.ota.OtaPackageEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.queue.QueueEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.relation.RelationEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.resource.ResourceEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.rule.RuleChainEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.settings.AdminSettingsEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.telemetry.TelemetryEdgeProcessor; @@ -139,6 +141,9 @@ public class EdgeContextComponent { @Autowired private QueueService queueService; + @Autowired + private ResourceService resourceService; + @Autowired private AlarmEdgeProcessor alarmProcessor; @@ -199,6 +204,9 @@ public class EdgeContextComponent { @Autowired private TenantProfileEdgeProcessor tenantProfileEdgeProcessor; + @Autowired + private ResourceEdgeProcessor resourceEdgeProcessor; + @Autowired private EdgeMsgConstructor edgeMsgConstructor; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index cfb92d0bfd..cf3f325260 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -63,6 +63,7 @@ import org.thingsboard.server.gen.edge.v1.RelationRequestMsg; import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg; import org.thingsboard.server.gen.edge.v1.RequestMsg; import org.thingsboard.server.gen.edge.v1.RequestMsgType; +import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.v1.SyncCompletedMsg; @@ -645,6 +646,8 @@ public final class EdgeGrpcSession implements Closeable { return ctx.getAdminSettingsProcessor().convertAdminSettingsEventToDownlink(edgeEvent); case OTA_PACKAGE: return ctx.getOtaPackageEdgeProcessor().convertOtaPackageEventToDownlink(edgeEvent); + case TB_RESOURCE: + return ctx.getResourceEdgeProcessor().convertResourceEventToDownlink(edgeEvent); case QUEUE: return ctx.getQueueEdgeProcessor().convertQueueEventToDownlink(edgeEvent); case TENANT: @@ -710,6 +713,11 @@ public final class EdgeGrpcSession implements Closeable { result.add(ctx.getDashboardProcessor().processDashboardMsgFromEdge(edge.getTenantId(), edge, dashboardUpdateMsg)); } } + if (uplinkMsg.getResourceUpdateMsgCount() > 0) { + for (ResourceUpdateMsg resourceUpdateMsg : uplinkMsg.getResourceUpdateMsgList()) { + result.add(ctx.getResourceEdgeProcessor().processResourceMsgFromEdge(edge.getTenantId(), resourceUpdateMsg)); + } + } if (uplinkMsg.getRuleChainMetadataRequestMsgCount() > 0) { for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) { result.add(ctx.getEdgeRequestsService().processRuleChainMetadataRequestMsg(edge.getTenantId(), edge, ruleChainMetadataRequestMsg)); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java index a39dae3b54..692a46ff44 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java @@ -33,10 +33,12 @@ import org.thingsboard.server.service.edge.rpc.fetch.EntityViewsEdgeEventFetcher import org.thingsboard.server.service.edge.rpc.fetch.OtaPackagesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.QueuesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.SystemResourcesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetTypesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.TenantEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.TenantResourcesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetTypesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetsBundlesEdgeEventFetcher; @@ -77,6 +79,8 @@ public class EdgeSyncCursor { fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService())); + fetchers.add(new SystemResourcesEdgeEventFetcher(ctx.getResourceService())); + fetchers.add(new TenantResourcesEdgeEventFetcher(ctx.getResourceService())); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/ResourceMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/ResourceMsgConstructor.java new file mode 100644 index 0000000000..27b7dc3978 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/ResourceMsgConstructor.java @@ -0,0 +1,57 @@ +/** + * Copyright © 2016-2023 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 org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.id.TbResourceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.queue.util.TbCoreComponent; + +@Component +@TbCoreComponent +public class ResourceMsgConstructor { + + public ResourceUpdateMsg constructResourceUpdatedMsg(UpdateMsgType msgType, TbResource tbResource) { + ResourceUpdateMsg.Builder builder = ResourceUpdateMsg.newBuilder() + .setMsgType(msgType) + .setIdMSB(tbResource.getId().getId().getMostSignificantBits()) + .setIdLSB(tbResource.getId().getId().getLeastSignificantBits()) + .setTitle(tbResource.getTitle()) + .setResourceKey(tbResource.getResourceKey()) + .setResourceType(tbResource.getResourceType().name()) + .setFileName(tbResource.getFileName()); + if (tbResource.getData() != null) { + builder.setData(tbResource.getData()); + } + if (tbResource.getEtag() != null) { + builder.setEtag(tbResource.getEtag()); + } + if (tbResource.getTenantId().equals(TenantId.SYS_TENANT_ID)) { + builder.setIsSystem(true); + } + return builder.build(); + } + + public ResourceUpdateMsg constructResourceDeleteMsg(TbResourceId tbResourceId) { + return ResourceUpdateMsg.newBuilder() + .setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) + .setIdMSB(tbResourceId.getId().getMostSignificantBits()) + .setIdLSB(tbResourceId.getId().getLeastSignificantBits()).build(); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseResourceEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseResourceEdgeEventFetcher.java new file mode 100644 index 0000000000..5e0ea661e0 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseResourceEdgeEventFetcher.java @@ -0,0 +1,49 @@ +/** + * Copyright © 2016-2023 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.TbResource; +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.resource.ResourceService; + +@Slf4j +@AllArgsConstructor +public abstract class BaseResourceEdgeEventFetcher extends BasePageableEdgeEventFetcher { + + protected final ResourceService resourceService; + + @Override + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return findTenantResources(tenantId, pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, TbResource tbResource) { + return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.TB_RESOURCE, + EdgeEventActionType.ADDED, tbResource.getId(), null); + } + + protected abstract PageData findTenantResources(TenantId tenantId, PageLink pageLink); +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/SystemResourcesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/SystemResourcesEdgeEventFetcher.java new file mode 100644 index 0000000000..42700ded3e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/SystemResourcesEdgeEventFetcher.java @@ -0,0 +1,34 @@ +/** + * Copyright © 2016-2023 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 org.thingsboard.server.common.data.TbResource; +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.resource.ResourceService; + +public class SystemResourcesEdgeEventFetcher extends BaseResourceEdgeEventFetcher { + + public SystemResourcesEdgeEventFetcher(ResourceService resourceService) { + super(resourceService); + } + + @Override + protected PageData findTenantResources(TenantId tenantId, PageLink pageLink) { + return resourceService.findAllTenantResources(TenantId.SYS_TENANT_ID, pageLink); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantResourcesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantResourcesEdgeEventFetcher.java new file mode 100644 index 0000000000..0992565285 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantResourcesEdgeEventFetcher.java @@ -0,0 +1,36 @@ +/** + * Copyright © 2016-2023 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.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.TbResource; +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.resource.ResourceService; + +@Slf4j +public class TenantResourcesEdgeEventFetcher extends BaseResourceEdgeEventFetcher { + + public TenantResourcesEdgeEventFetcher(ResourceService resourceService) { + super(resourceService); + } + + @Override + protected PageData findTenantResources(TenantId tenantId, PageLink pageLink) { + return resourceService.findAllTenantResources(tenantId, pageLink); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java index ffddb65315..8f69a904e8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java @@ -28,6 +28,7 @@ public class TenantWidgetsBundlesEdgeEventFetcher extends BaseWidgetsBundlesEdge public TenantWidgetsBundlesEdgeEventFetcher(WidgetsBundleService widgetsBundleService) { super(widgetsBundleService); } + @Override protected PageData findWidgetsBundles(TenantId tenantId, PageLink pageLink) { return widgetsBundleService.findTenantWidgetsBundlesByTenantId(tenantId, pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index 328a64efe2..99375dc747 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.edge.Edge; @@ -72,6 +73,7 @@ import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.tenant.TenantProfileService; @@ -101,6 +103,7 @@ import org.thingsboard.server.service.edge.rpc.constructor.OtaPackageMsgConstruc import org.thingsboard.server.service.edge.rpc.constructor.QueueMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RelationMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RuleChainMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.ResourceMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.TenantMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.TenantProfileMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.UserMsgConstructor; @@ -211,6 +214,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected PartitionService partitionService; + @Autowired + protected ResourceService resourceService; + @Autowired @Lazy protected TbQueueProducerProvider producerProvider; @@ -233,6 +239,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected DataValidator entityViewValidator; + @Autowired + protected DataValidator resourceValidator; + @Autowired protected EdgeMsgConstructor edgeMsgConstructor; @@ -293,6 +302,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected QueueMsgConstructor queueMsgConstructor; + @Autowired + protected ResourceMsgConstructor resourceMsgConstructor; + @Autowired protected EdgeSynchronizationManager edgeSynchronizationManager; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/BaseResourceProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/BaseResourceProcessor.java new file mode 100644 index 0000000000..13ca012f3c --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/BaseResourceProcessor.java @@ -0,0 +1,61 @@ +/** + * Copyright © 2016-2023 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.resource; + +import com.datastax.oss.driver.api.core.uuid.Uuids; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.ResourceType; +import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.TbResourceInfo; +import org.thingsboard.server.common.data.id.TbResourceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; +import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; + +@Slf4j +public abstract class BaseResourceProcessor extends BaseEdgeProcessor { + + protected void saveOrUpdateTbResource(TenantId tenantId, TbResourceId tbResourceId, ResourceUpdateMsg resourceUpdateMsg) { + try { + boolean created = false; + TbResource resource = resourceService.findResourceById(tenantId, tbResourceId); + if (resource == null) { + resource = new TbResource(); + if (resourceUpdateMsg.getIsSystem()) { + resource.setTenantId(TenantId.SYS_TENANT_ID); + } else { + resource.setTenantId(tenantId); + } + resource.setCreatedTime(Uuids.unixTimestamp(tbResourceId.getId())); + created = true; + } + resource.setTitle(resourceUpdateMsg.getTitle()); + resource.setResourceKey(resourceUpdateMsg.getResourceKey()); + resource.setResourceType(ResourceType.valueOf(resourceUpdateMsg.getResourceType())); + resource.setFileName(resourceUpdateMsg.getFileName()); + resource.setData(resourceUpdateMsg.hasData() ? resourceUpdateMsg.getData() : null); + resource.setEtag(resourceUpdateMsg.hasEtag() ? resourceUpdateMsg.getEtag() : null); + resourceValidator.validate(resource, TbResourceInfo::getTenantId); + if (created) { + resource.setId(tbResourceId); + } + resourceService.saveResource(resource, false); + } catch (Exception e) { + log.error("[{}] Failed to process resource update msg [{}]", tenantId, resourceUpdateMsg, e); + throw e; + } + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java new file mode 100644 index 0000000000..ebc9a5440e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java @@ -0,0 +1,100 @@ +/** + * Copyright © 2016-2023 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.resource; + +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.EdgeUtils; +import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.id.TbResourceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.gen.edge.v1.DownlinkMsg; +import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.queue.util.TbCoreComponent; + +import java.util.UUID; + +@Component +@Slf4j +@TbCoreComponent +public class ResourceEdgeProcessor extends BaseResourceProcessor { + + public ListenableFuture processResourceMsgFromEdge(TenantId tenantId, ResourceUpdateMsg resourceUpdateMsg) { + TbResourceId tbResourceId = new TbResourceId(new UUID(resourceUpdateMsg.getIdMSB(), resourceUpdateMsg.getIdLSB())); + try { + edgeSynchronizationManager.getSync().set(true); + + switch (resourceUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE: + case ENTITY_UPDATED_RPC_MESSAGE: + super.saveOrUpdateTbResource(tenantId, tbResourceId, resourceUpdateMsg); + break; + case ENTITY_DELETED_RPC_MESSAGE: + TbResource tbResourceToDelete = resourceService.findResourceById(tenantId, tbResourceId); + if (tbResourceToDelete != null) { + resourceService.deleteResource(tenantId, tbResourceId); + } + break; + case UNRECOGNIZED: + return handleUnsupportedMsgType(resourceUpdateMsg.getMsgType()); + } + } catch (DataValidationException e) { + if (e.getMessage().contains("files size limit is exhausted")) { + log.warn("[{}] Resource data size has been exhausted {}", tenantId, resourceUpdateMsg, e); + return Futures.immediateFuture(null); + } else { + return Futures.immediateFailedFuture(e); + } + } finally { + edgeSynchronizationManager.getSync().remove(); + } + return Futures.immediateFuture(null); + } + + public DownlinkMsg convertResourceEventToDownlink(EdgeEvent edgeEvent) { + TbResourceId tbResourceId = new TbResourceId(edgeEvent.getEntityId()); + DownlinkMsg downlinkMsg = null; + switch (edgeEvent.getAction()) { + case ADDED: + case UPDATED: + TbResource tbResource = resourceService.findResourceById(edgeEvent.getTenantId(), tbResourceId); + if (tbResource != null) { + UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); + ResourceUpdateMsg resourceUpdateMsg = + resourceMsgConstructor.constructResourceUpdatedMsg(msgType, tbResource); + downlinkMsg = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addResourceUpdateMsg(resourceUpdateMsg) + .build(); + } + break; + case DELETED: + ResourceUpdateMsg resourceUpdateMsg = + resourceMsgConstructor.constructResourceDeleteMsg(tbResourceId); + downlinkMsg = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addResourceUpdateMsg(resourceUpdateMsg) + .build(); + break; + } + return downlinkMsg; + } +} diff --git a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java index 410348762f..a54f0e2a5d 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.edge; -import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.google.protobuf.InvalidProtocolBufferException; @@ -393,7 +392,7 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { Assert.assertEquals(expectedRuleChainUUID, ruleChainUUID); } - private void validateAdminSettings() throws JsonProcessingException { + private void validateAdminSettings() { List adminSettingsUpdateMsgs = edgeImitator.findAllMessagesByType(AdminSettingsUpdateMsg.class); Assert.assertEquals(4, adminSettingsUpdateMsgs.size()); diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessorTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessorTest.java index d8c8041694..fcad4a4500 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessorTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessorTest.java @@ -40,5 +40,4 @@ class AssetEdgeProcessorTest extends AbstractAssetProcessorTest { verify(downlinkMsg, expectedDashboardIdMSB, expectedDashboardIdLSB, expectedRuleChainIdMSB, expectedRuleChainIdLSB); } - -} \ No newline at end of file +} diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetProfileEdgeProcessorTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetProfileEdgeProcessorTest.java index e0a6155664..bafad59050 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetProfileEdgeProcessorTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetProfileEdgeProcessorTest.java @@ -40,4 +40,4 @@ class AssetProfileEdgeProcessorTest extends AbstractAssetProcessorTest{ verify(downlinkMsg, expectedDashboardIdMSB, expectedDashboardIdLSB, expectedRuleChainIdMSB, expectedRuleChainIdLSB); } -} \ No newline at end of file +} diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessorTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessorTest.java index f64bd898d1..e40ea364be 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessorTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessorTest.java @@ -39,5 +39,4 @@ class DeviceEdgeProcessorTest extends AbstractDeviceProcessorTest { verify(downlinkMsg, expectedDashboardIdMSB, expectedDashboardIdLSB, expectedRuleChainIdMSB, expectedRuleChainIdLSB); } - -} \ No newline at end of file +} diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceProfileEdgeProcessorTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceProfileEdgeProcessorTest.java index 35d44852e9..c2e9dca077 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceProfileEdgeProcessorTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceProfileEdgeProcessorTest.java @@ -41,5 +41,4 @@ class DeviceProfileEdgeProcessorTest extends AbstractDeviceProcessorTest { verify(downlinkMsg, expectedDashboardIdMSB, expectedDashboardIdLSB, expectedRuleChainIdMSB, expectedRuleChainIdLSB); } - -} \ No newline at end of file +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java index 6f0e362209..5bc45ebc29 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/resource/ResourceService.java @@ -32,12 +32,16 @@ public interface ResourceService extends EntityDaoService { TbResource saveResource(TbResource resource); + TbResource saveResource(TbResource resource, boolean doValidate); + TbResource getResource(TenantId tenantId, ResourceType resourceType, String resourceId); TbResource findResourceById(TenantId tenantId, TbResourceId resourceId); TbResourceInfo findResourceInfoById(TenantId tenantId, TbResourceId resourceId); + PageData findAllTenantResources(TenantId tenantId, PageLink pageLink); + ListenableFuture findResourceInfoByIdAsync(TenantId tenantId, TbResourceId resourceId); PageData findAllTenantResourcesByTenantId(TbResourceInfoFilter filter, PageLink pageLink); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java index 1397deb5d7..5ecd8eab27 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java @@ -39,7 +39,8 @@ public enum EdgeEventType { WIDGET_TYPE(true, EntityType.WIDGET_TYPE), ADMIN_SETTINGS(true, null), OTA_PACKAGE(true, EntityType.OTA_PACKAGE), - QUEUE(true, EntityType.QUEUE); + QUEUE(true, EntityType.QUEUE), + TB_RESOURCE(true, EntityType.TB_RESOURCE); private final boolean allEdgesRelated; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java index 0cdf3ad1eb..cc66919c0f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java @@ -139,6 +139,8 @@ public class EntityIdFactory { return new EdgeId(uuid); case QUEUE: return new QueueId(uuid); + case TB_RESOURCE: + return new TbResourceId(uuid); } throw new IllegalArgumentException("EdgeEventType " + edgeEventType + " is not supported!"); } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 50d5ad3c7f..ff117f864d 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -420,6 +420,19 @@ message TenantProfileUpdateMsg { bytes profileDataBytes = 8; } +message ResourceUpdateMsg { + UpdateMsgType msgType = 1; + int64 idMSB = 2; + int64 idLSB = 3; + string title = 4; + string resourceType = 5; + string resourceKey = 6; + string fileName = 7; + optional string data = 8; + optional string etag = 9; + bool isSystem = 10; +} + message RuleChainMetadataRequestMsg { int64 ruleChainIdMSB = 1; int64 ruleChainIdLSB = 2; @@ -571,6 +584,7 @@ message UplinkMsg { repeated EntityViewUpdateMsg entityViewUpdateMsg = 18; repeated AssetProfileUpdateMsg assetProfileUpdateMsg = 19; repeated DeviceProfileUpdateMsg deviceProfileUpdateMsg = 20; + repeated ResourceUpdateMsg resourceUpdateMsg = 21; } message UplinkResponseMsg { @@ -613,5 +627,6 @@ message DownlinkMsg { EdgeConfiguration edgeConfiguration = 25; repeated TenantUpdateMsg tenantUpdateMsg = 26; repeated TenantProfileUpdateMsg tenantProfileUpdateMsg = 27; + repeated ResourceUpdateMsg resourceUpdateMsg = 28; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java index 7697217b6b..76ec3766d7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java @@ -21,10 +21,9 @@ import lombok.extern.slf4j.Slf4j; import org.hibernate.exception.ConstraintViolationException; import org.springframework.stereotype.Service; import org.springframework.transaction.event.TransactionalEventListener; -import org.thingsboard.server.cache.device.DeviceCacheKey; +import org.thingsboard.server.cache.resourceInfo.ResourceInfoCacheKey; import org.thingsboard.server.cache.resourceInfo.ResourceInfoEvictEvent; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.cache.resourceInfo.ResourceInfoCacheKey; import org.thingsboard.server.common.data.ResourceType; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.TbResourceInfo; @@ -36,6 +35,8 @@ 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.entity.AbstractCachedEntityService; +import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; +import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; @@ -59,10 +60,23 @@ public class BaseResourceService extends AbstractCachedEntityService findAllTenantResources(TenantId tenantId, PageLink pageLink) { + log.trace("Executing findAllTenantResources [{}][{}]", tenantId, pageLink); + validateId(tenantId, INCORRECT_TENANT_ID + tenantId); + return resourceDao.findAllByTenantId(tenantId, pageLink); + } + @Override public PageData findTenantResourcesByResourceTypeAndPageLink(TenantId tenantId, ResourceType resourceType, PageLink pageLink) { log.trace("Executing findTenantResourcesByResourceTypeAndPageLink [{}][{}][{}]", tenantId, resourceType, pageLink); From c18264974581074bd9005511abf0d4bd889927d6 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Tue, 12 Sep 2023 10:40:00 +0300 Subject: [PATCH 2/5] Add test to check resource creation from edge and fix existing ones --- .../controller/TbResourceControllerTest.java | 43 ++++++------ .../server/edge/ResourceEdgeTest.java | 66 +++++++++++++++++++ .../rpc/processor/BaseEdgeProcessorTest.java | 12 ++++ .../dao/resource/BaseResourceService.java | 2 +- 4 files changed, 100 insertions(+), 23 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java diff --git a/application/src/test/java/org/thingsboard/server/controller/TbResourceControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/TbResourceControllerTest.java index 8725f43be1..5820ae8afc 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TbResourceControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/TbResourceControllerTest.java @@ -20,7 +20,6 @@ import com.fasterxml.jackson.databind.JsonNode; import org.junit.After; import org.junit.Assert; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.mockito.Mockito; import org.springframework.http.HttpHeaders; @@ -39,7 +38,6 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.widget.WidgetTypeDetails; -import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DaoSqlTest; @@ -105,9 +103,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { TbResource savedResource = save(resource); - testNotifyEntityOneTimeMsgToEdgeServiceNever(savedResource, savedResource.getId(), savedResource.getId(), + testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(savedResource, savedResource.getId(), savedResource.getId(), savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), - ActionType.ADDED); + ActionType.ADDED, ActionType.ADDED); Assert.assertNotNull(savedResource); Assert.assertNotNull(savedResource.getId()); @@ -125,9 +123,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { TbResource foundResource = doGet("/api/resource/" + savedResource.getId().getId().toString(), TbResource.class); Assert.assertEquals(foundResource.getTitle(), savedResource.getTitle()); - testNotifyEntityOneTimeMsgToEdgeServiceNever(foundResource, foundResource.getId(), foundResource.getId(), + testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(foundResource, foundResource.getId(), foundResource.getId(), savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), - ActionType.UPDATED); + ActionType.UPDATED, ActionType.UPDATED); } @Test @@ -157,7 +155,7 @@ public class TbResourceControllerTest extends AbstractControllerTest { resource.setFileName(DEFAULT_FILE_NAME); resource.setData(TEST_DATA); - TbResource savedResource = save(resource); + TbResource savedResource = save(resource); loginDifferentTenant(); @@ -209,9 +207,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { .andExpect(status().isOk()); - testNotifyEntityOneTimeMsgToEdgeServiceNever(savedResource, savedResource.getId(), savedResource.getId(), + testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(savedResource, savedResource.getId(), savedResource.getId(), savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), - ActionType.DELETED, resourceIdStr); + ActionType.DELETED, ActionType.DELETED, resourceIdStr); doGet("/api/resource/" + savedResource.getId().getId().toString()) .andExpect(status().isNotFound()) @@ -240,7 +238,7 @@ public class TbResourceControllerTest extends AbstractControllerTest { doDelete("/api/resource/" + resourceIdStr) .andExpect(status().isBadRequest()) .andExpect(statusReason(containsString("Following widget types uses current resource: [" - + widgetType .getName()+ "]"))); + + widgetType.getName() + "]"))); } @Test @@ -271,9 +269,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { } } while (pageData.hasNext()); - testNotifyManyEntityManyTimeMsgToEdgeServiceNever(new TbResource(), new TbResource(), + testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(new TbResource(), new TbResource(), savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), - ActionType.ADDED, cntEntity); + ActionType.ADDED, cntEntity, cntEntity, cntEntity); Collections.sort(resources, idComparator); Collections.sort(loadedResources, idComparator); @@ -319,9 +317,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { } } while (pageData.hasNext()); - testNotifyManyEntityManyTimeMsgToEdgeServiceNever(new TbResource(), new TbResource(), - savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), - ActionType.ADDED, jksCntEntity + lwm2mCntEntity); + testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(new TbResource(), new TbResource(), + savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), ActionType.ADDED, + jksCntEntity + lwm2mCntEntity, jksCntEntity + lwm2mCntEntity, jksCntEntity + lwm2mCntEntity); Collections.sort(resources, idComparator); Collections.sort(loadedResources, idComparator); @@ -368,9 +366,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { .andExpect(status().isOk()); } - testNotifyManyEntityManyTimeMsgToEdgeServiceNeverAdditionalInfoAny(new TbResource(), new TbResource(), + testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAnyAdditionalInfoAny(new TbResource(), new TbResource(), resources.get(0).getTenantId(), null, null, SYS_ADMIN_EMAIL, - ActionType.DELETED, cntEntity, 1); + ActionType.DELETED, ActionType.DELETED, cntEntity, cntEntity, 1); pageLink = new PageLink(27); loadedResources.clear(); @@ -441,9 +439,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { .andExpect(status().isOk()); } - testNotifyManyEntityManyTimeMsgToEdgeServiceNeverAdditionalInfoAny(new TbResource(), new TbResource(), + testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAnyAdditionalInfoAny(new TbResource(), new TbResource(), jksResources.get(0).getTenantId(), null, null, SYS_ADMIN_EMAIL, - ActionType.DELETED, cntEntity, 1); + ActionType.DELETED, ActionType.DELETED, cntEntity, cntEntity, 1); pageLink = new PageLink(27); loadedResources.clear(); @@ -535,9 +533,9 @@ public class TbResourceControllerTest extends AbstractControllerTest { TbResource savedResource = save(resource); - testNotifyEntityOneTimeMsgToEdgeServiceNever(savedResource, savedResource.getId(), savedResource.getId(), + testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(savedResource, savedResource.getId(), savedResource.getId(), savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), - ActionType.ADDED); + ActionType.ADDED, ActionType.ADDED); ResultActions resultActions = doGet("/api/resource/js/" + savedResource.getId().getId().toString() + "/download") .andExpect(status().isOk()); @@ -617,6 +615,7 @@ public class TbResourceControllerTest extends AbstractControllerTest { } private TbResource save(TbResource tbResource) throws Exception { - return doPostWithTypedResponse("/api/resource", tbResource, new TypeReference<>(){}); + return doPostWithTypedResponse("/api/resource", tbResource, new TypeReference<>() { + }); } } diff --git a/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java new file mode 100644 index 0000000000..554f4ca7cf --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java @@ -0,0 +1,66 @@ +/** + * Copyright © 2016-2023 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.edge; + +import com.datastax.oss.driver.api.core.uuid.Uuids; +import org.junit.Assert; +import org.junit.Test; +import org.thingsboard.server.common.data.ResourceType; +import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.edge.v1.UplinkMsg; +import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; + +import java.util.UUID; + +@DaoSqlTest +public class ResourceEdgeTest extends AbstractEdgeTest { + + @Test + public void testSendResourceToCloud() throws Exception { + UUID uuid = Uuids.timeBased(); + + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + ResourceUpdateMsg.Builder resourceUpdateMsgBuilder = ResourceUpdateMsg.newBuilder(); + resourceUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits()); + resourceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); + resourceUpdateMsgBuilder.setTitle("Edge Test Resource"); + resourceUpdateMsgBuilder.setResourceType(ResourceType.JKS.name()); + resourceUpdateMsgBuilder.setResourceKey("EdgeResource.jks"); + resourceUpdateMsgBuilder.setFileName("EdgeResource.jks"); + resourceUpdateMsgBuilder.setData("Data"); + resourceUpdateMsgBuilder.setIsSystem(false); + resourceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE); + testAutoGeneratedCodeByProtobuf(resourceUpdateMsgBuilder); + uplinkMsgBuilder.addResourceUpdateMsg(resourceUpdateMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + + Assert.assertTrue(edgeImitator.waitForResponses()); + + UplinkResponseMsg latestResponseMsg = edgeImitator.getLatestResponseMsg(); + Assert.assertTrue(latestResponseMsg.getSuccess()); + + TbResource tbResource = doGet("/api/resource/" + uuid, TbResource.class); + Assert.assertNotNull(tbResource); + Assert.assertEquals("Edge Test Resource", tbResource.getName()); + } +} diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessorTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessorTest.java index 1c4b481302..dc59d4f7e9 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessorTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessorTest.java @@ -24,6 +24,7 @@ import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.edge.EdgeEvent; @@ -47,6 +48,7 @@ import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.tenant.TenantProfileService; @@ -72,6 +74,7 @@ import org.thingsboard.server.service.edge.rpc.constructor.EntityViewMsgConstruc import org.thingsboard.server.service.edge.rpc.constructor.OtaPackageMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.QueueMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RelationMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.ResourceMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RuleChainMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.TenantMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.TenantProfileMsgConstructor; @@ -174,6 +177,9 @@ public abstract class BaseEdgeProcessorTest { @MockBean protected PartitionService partitionService; + @MockBean + protected ResourceService resourceService; + @MockBean @Lazy protected TbQueueProducerProvider producerProvider; @@ -196,6 +202,9 @@ public abstract class BaseEdgeProcessorTest { @MockBean protected DataValidator entityViewValidator; + @MockBean + protected DataValidator resourceValidator; + @MockBean protected EdgeMsgConstructor edgeMsgConstructor; @@ -256,6 +265,9 @@ public abstract class BaseEdgeProcessorTest { @MockBean protected QueueMsgConstructor queueMsgConstructor; + @MockBean + protected ResourceMsgConstructor resourceMsgConstructor; + @MockBean protected EdgeSynchronizationManager edgeSynchronizationManager; diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java index 76ec3766d7..2181318eed 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java @@ -76,7 +76,7 @@ public class BaseResourceService extends AbstractCachedEntityService Date: Fri, 27 Oct 2023 13:06:06 +0300 Subject: [PATCH 3/5] Add entity notification for resources --- .../server/service/edge/DefaultEdgeNotificationService.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index f213169299..87e5adeb99 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -225,6 +225,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { case TENANT_PROFILE: tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break; + case TB_RESOURCE: + resourceEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; default: log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); } From d92c016ad9fdb2ba56f6050ebe94be706ef4a1c8 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 27 Oct 2023 17:46:41 +0300 Subject: [PATCH 4/5] Added ResourceEdgeTest.testResources_create_update_delete --- .../server/edge/ResourceEdgeTest.java | 63 ++++++++++++++++++- .../server/edge/imitator/EdgeImitator.java | 6 ++ 2 files changed, 66 insertions(+), 3 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java index 554f4ca7cf..6a3001f292 100644 --- a/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/ResourceEdgeTest.java @@ -16,9 +16,11 @@ package org.thingsboard.server.edge; import com.datastax.oss.driver.api.core.uuid.Uuids; +import com.google.protobuf.AbstractMessage; import org.junit.Assert; import org.junit.Test; import org.thingsboard.server.common.data.ResourceType; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; @@ -28,9 +30,64 @@ import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import java.util.UUID; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; + @DaoSqlTest public class ResourceEdgeTest extends AbstractEdgeTest { + private static final String TEST_DATA = "77u/PD94bWwgdmVyc2lvbj0iMS4wIiBlbmNvZGluZz0iVVRGLTgiPz4KPCEtLQpGSUxFIElORk9STUFUSU9OCgpPTUEgUGVybWFuZW50IERvY3VtZW50CiAgIEZpbGU6IE9NQS1TVVAtTHdNMk1fQmluYXJ5QXBwRGF0YUNvbnRhaW5lci1WMV8wXzEtMjAxOTAyMjEtQQogICBUeXBlOiB4bWwKClB1YmxpYyBSZWFjaGFibGUgSW5mb3JtYXRpb24KICAgUGF0aDogaHR0cDovL3d3dy5vcGVubW9iaWxlYWxsaWFuY2Uub3JnL3RlY2gvcHJvZmlsZXMKICAgTmFtZTogTHdNMk1fQmluYXJ5QXBwRGF0YUNvbnRhaW5lci12MV8wXzEueG1sCgpOT1JNQVRJVkUgSU5GT1JNQVRJT04KCiAgSW5mb3JtYXRpb24gYWJvdXQgdGhpcyBmaWxlIGNhbiBiZSBmb3VuZCBpbiB0aGUgbGF0ZXN0IHJldmlzaW9uIG9mCgogIE9NQS1UUy1MV00yTV9CaW5hcnlBcHBEYXRhQ29udGFpbmVyLVYxXzBfMQoKICBUaGlzIGlzIGF2YWlsYWJsZSBhdCBodHRwOi8vd3d3Lm9wZW5tb2JpbGVhbGxpYW5jZS5vcmcvCgogIFNlbmQgY29tbWVudHMgdG8gaHR0cHM6Ly9naXRodWIuY29tL09wZW5Nb2JpbGVBbGxpYW5jZS9PTUFfTHdNMk1fZm9yX0RldmVsb3BlcnMvaXNzdWVzCgpDSEFOR0UgSElTVE9SWQoKMTUwNjIwMTggU3RhdHVzIGNoYW5nZWQgdG8gQXBwcm92ZWQgYnkgRE0sIERvYyBSZWYgIyBPTUEtRE0mU0UtMjAxOC0wMDYxLUlOUF9MV00yTV9BUFBEQVRBX1YxXzBfRVJQX2Zvcl9maW5hbF9BcHByb3ZhbAoyMTAyMjAxOSBTdGF0dXMgY2hhbmdlZCB0byBBcHByb3ZlZCBieSBJUFNPLCBEb2MgUmVmICMgT01BLUlQU08tMjAxOS0wMDI1LUlOUF9Md00yTV9PYmplY3RfQXBwX0RhdGFfQ29udGFpbmVyXzFfMF8xX2Zvcl9GaW5hbF9BcHByb3ZhbAoKTEVHQUwgRElTQ0xBSU1FUgoKQ29weXJpZ2h0IDIwMTkgT3BlbiBNb2JpbGUgQWxsaWFuY2UuCgpSZWRpc3RyaWJ1dGlvbiBhbmQgdXNlIGluIHNvdXJjZSBhbmQgYmluYXJ5IGZvcm1zLCB3aXRoIG9yIHdpdGhvdXQKbW9kaWZpY2F0aW9uLCBhcmUgcGVybWl0dGVkIHByb3ZpZGVkIHRoYXQgdGhlIGZvbGxvd2luZyBjb25kaXRpb25zCmFyZSBtZXQ6CgoxLiBSZWRpc3RyaWJ1dGlvbnMgb2Ygc291cmNlIGNvZGUgbXVzdCByZXRhaW4gdGhlIGFib3ZlIGNvcHlyaWdodApub3RpY2UsIHRoaXMgbGlzdCBvZiBjb25kaXRpb25zIGFuZCB0aGUgZm9sbG93aW5nIGRpc2NsYWltZXIuCjIuIFJlZGlzdHJpYnV0aW9ucyBpbiBiaW5hcnkgZm9ybSBtdXN0IHJlcHJvZHVjZSB0aGUgYWJvdmUgY29weXJpZ2h0Cm5vdGljZSwgdGhpcyBsaXN0IG9mIGNvbmRpdGlvbnMgYW5kIHRoZSBmb2xsb3dpbmcgZGlzY2xhaW1lciBpbiB0aGUKZG9jdW1lbnRhdGlvbiBhbmQvb3Igb3RoZXIgbWF0ZXJpYWxzIHByb3ZpZGVkIHdpdGggdGhlIGRpc3RyaWJ1dGlvbi4KMy4gTmVpdGhlciB0aGUgbmFtZSBvZiB0aGUgY29weXJpZ2h0IGhvbGRlciBub3IgdGhlIG5hbWVzIG9mIGl0cwpjb250cmlidXRvcnMgbWF5IGJlIHVzZWQgdG8gZW5kb3JzZSBvciBwcm9tb3RlIHByb2R1Y3RzIGRlcml2ZWQKZnJvbSB0aGlzIHNvZnR3YXJlIHdpdGhvdXQgc3BlY2lmaWMgcHJpb3Igd3JpdHRlbiBwZXJtaXNzaW9uLgoKVEhJUyBTT0ZUV0FSRSBJUyBQUk9WSURFRCBCWSBUSEUgQ09QWVJJR0hUIEhPTERFUlMgQU5EIENPTlRSSUJVVE9SUwoiQVMgSVMiIEFORCBBTlkgRVhQUkVTUyBPUiBJTVBMSUVEIFdBUlJBTlRJRVMsIElOQ0xVRElORywgQlVUIE5PVApMSU1JVEVEIFRPLCBUSEUgSU1QTElFRCBXQVJSQU5USUVTIE9GIE1FUkNIQU5UQUJJTElUWSBBTkQgRklUTkVTUwpGT1IgQSBQQVJUSUNVTEFSIFBVUlBPU0UgQVJFIERJU0NMQUlNRUQuIElOIE5PIEVWRU5UIFNIQUxMIFRIRQpDT1BZUklHSFQgSE9MREVSIE9SIENPTlRSSUJVVE9SUyBCRSBMSUFCTEUgRk9SIEFOWSBESVJFQ1QsIElORElSRUNULApJTkNJREVOVEFMLCBTUEVDSUFMLCBFWEVNUExBUlksIE9SIENPTlNFUVVFTlRJQUwgREFNQUdFUyAoSU5DTFVESU5HLApCVVQgTk9UIExJTUlURUQgVE8sIFBST0NVUkVNRU5UIE9GIFNVQlNUSVRVVEUgR09PRFMgT1IgU0VSVklDRVM7CkxPU1MgT0YgVVNFLCBEQVRBLCBPUiBQUk9GSVRTOyBPUiBCVVNJTkVTUyBJTlRFUlJVUFRJT04pIEhPV0VWRVIKQ0FVU0VEIEFORCBPTiBBTlkgVEhFT1JZIE9GIExJQUJJTElUWSwgV0hFVEhFUiBJTiBDT05UUkFDVCwgU1RSSUNUCkxJQUJJTElUWSwgT1IgVE9SVCAoSU5DTFVESU5HIE5FR0xJR0VOQ0UgT1IgT1RIRVJXSVNFKSBBUklTSU5HIElOCkFOWSBXQVkgT1VUIE9GIFRIRSBVU0UgT0YgVEhJUyBTT0ZUV0FSRSwgRVZFTiBJRiBBRFZJU0VEIE9GIFRIRQpQT1NTSUJJTElUWSBPRiBTVUNIIERBTUFHRS4KClRoZSBhYm92ZSBsaWNlbnNlIGlzIHVzZWQgYXMgYSBsaWNlbnNlIHVuZGVyIGNvcHlyaWdodCBvbmx5LiBQbGVhc2UKcmVmZXJlbmNlIHRoZSBPTUEgSVBSIFBvbGljeSBmb3IgcGF0ZW50IGxpY2Vuc2luZyB0ZXJtczoKaHR0cHM6Ly93d3cub21hc3BlY3dvcmtzLm9yZy9hYm91dC9pbnRlbGxlY3R1YWwtcHJvcGVydHktcmlnaHRzLwoKLS0+CjxMV00yTSB4bWxuczp4c2k9Imh0dHA6Ly93d3cudzMub3JnLzIwMDEvWE1MU2NoZW1hLWluc3RhbmNlIiB4c2k6bm9OYW1lc3BhY2VTY2hlbWFMb2NhdGlvbj0iaHR0cDovL29wZW5tb2JpbGVhbGxpYW5jZS5vcmcvdGVjaC9wcm9maWxlcy9MV00yTS54c2QiPgoJPE9iamVjdCBPYmplY3RUeXBlPSJNT0RlZmluaXRpb24iPgoJCTxOYW1lPkJpbmFyeUFwcERhdGFDb250YWluZXI8L05hbWU+CgkJPERlc2NyaXB0aW9uMT48IVtDREFUQVtUaGlzIEx3TTJNIE9iamVjdHMgcHJvdmlkZXMgdGhlIGFwcGxpY2F0aW9uIHNlcnZpY2UgZGF0YSByZWxhdGVkIHRvIGEgTHdNMk0gU2VydmVyLCBlZy4gV2F0ZXIgbWV0ZXIgZGF0YS4gClRoZXJlIGFyZSBzZXZlcmFsIG1ldGhvZHMgdG8gY3JlYXRlIGluc3RhbmNlIHRvIGluZGljYXRlIHRoZSBtZXNzYWdlIGRpcmVjdGlvbiBiYXNlZCBvbiB0aGUgbmVnb3RpYXRpb24gYmV0d2VlbiBBcHBsaWNhdGlvbiBhbmQgTHdNMk0uIFRoZSBDbGllbnQgYW5kIFNlcnZlciBzaG91bGQgbmVnb3RpYXRlIHRoZSBpbnN0YW5jZShzKSB1c2VkIHRvIGV4Y2hhbmdlIHRoZSBkYXRhLiBGb3IgZXhhbXBsZToKIC0gVXNpbmcgYSBzaW5nbGUgaW5zdGFuY2UgZm9yIGJvdGggZGlyZWN0aW9ucyBjb21tdW5pY2F0aW9uLCBmcm9tIENsaWVudCB0byBTZXJ2ZXIgYW5kIGZyb20gU2VydmVyIHRvIENsaWVudC4KIC0gVXNpbmcgYW4gaW5zdGFuY2UgZm9yIGNvbW11bmljYXRpb24gZnJvbSBDbGllbnQgdG8gU2VydmVyIGFuZCBhbm90aGVyIG9uZSBmb3IgY29tbXVuaWNhdGlvbiBmcm9tIFNlcnZlciB0byBDbGllbnQKIC0gVXNpbmcgc2V2ZXJhbCBpbnN0YW5jZXMKXV0+PC9EZXNjcmlwdGlvbjE+CgkJPE9iamVjdElEPjE5PC9PYmplY3RJRD4KCQk8T2JqZWN0VVJOPnVybjpvbWE6bHdtMm06b21hOjE5PC9PYmplY3RVUk4+CgkJPExXTTJNVmVyc2lvbj4xLjA8L0xXTTJNVmVyc2lvbj4KCQk8T2JqZWN0VmVyc2lvbj4xLjA8L09iamVjdFZlcnNpb24+CgkJPE11bHRpcGxlSW5zdGFuY2VzPk11bHRpcGxlPC9NdWx0aXBsZUluc3RhbmNlcz4KCQk8TWFuZGF0b3J5Pk9wdGlvbmFsPC9NYW5kYXRvcnk+CgkJPFJlc291cmNlcz4KCQkJPEl0ZW0gSUQ9IjAiPjxOYW1lPkRhdGE8L05hbWU+CgkJCQk8T3BlcmF0aW9ucz5SVzwvT3BlcmF0aW9ucz4KCQkJCTxNdWx0aXBsZUluc3RhbmNlcz5NdWx0aXBsZTwvTXVsdGlwbGVJbnN0YW5jZXM+CgkJCQk8TWFuZGF0b3J5Pk1hbmRhdG9yeTwvTWFuZGF0b3J5PgoJCQkJPFR5cGU+T3BhcXVlPC9UeXBlPgoJCQkJPFJhbmdlRW51bWVyYXRpb24gLz4KCQkJCTxVbml0cyAvPgoJCQkJPERlc2NyaXB0aW9uPjwhW0NEQVRBW0luZGljYXRlcyB0aGUgYXBwbGljYXRpb24gZGF0YSBjb250ZW50Ll1dPjwvRGVzY3JpcHRpb24+CgkJCTwvSXRlbT4KCQkJPEl0ZW0gSUQ9IjEiPjxOYW1lPkRhdGEgUHJpb3JpdHk8L05hbWU+CgkJCQk8T3BlcmF0aW9ucz5SVzwvT3BlcmF0aW9ucz4KCQkJCTxNdWx0aXBsZUluc3RhbmNlcz5TaW5nbGU8L011bHRpcGxlSW5zdGFuY2VzPgoJCQkJPE1hbmRhdG9yeT5PcHRpb25hbDwvTWFuZGF0b3J5PgoJCQkJPFR5cGU+SW50ZWdlcjwvVHlwZT4KCQkJCTxSYW5nZUVudW1lcmF0aW9uPjEgYnl0ZXM8L1JhbmdlRW51bWVyYXRpb24+CgkJCQk8VW5pdHMgLz4KCQkJCTxEZXNjcmlwdGlvbj48IVtDREFUQVtJbmRpY2F0ZXMgdGhlIEFwcGxpY2F0aW9uIGRhdGEgcHJpb3JpdHk6CjA6SW1tZWRpYXRlCjE6QmVzdEVmZm9ydAoyOkxhdGVzdAozLTEwMDogUmVzZXJ2ZWQgZm9yIGZ1dHVyZSB1c2UuCjEwMS0yNTQ6IFByb3ByaWV0YXJ5IG1vZGUuXV0+PC9EZXNjcmlwdGlvbj4KCQkJPC9JdGVtPgoJCQk8SXRlbSBJRD0iMiI+PE5hbWU+RGF0YSBDcmVhdGlvbiBUaW1lPC9OYW1lPgoJCQkJPE9wZXJhdGlvbnM+Ulc8L09wZXJhdGlvbnM+CgkJCQk8TXVsdGlwbGVJbnN0YW5jZXM+U2luZ2xlPC9NdWx0aXBsZUluc3RhbmNlcz4KCQkJCTxNYW5kYXRvcnk+T3B0aW9uYWw8L01hbmRhdG9yeT4KCQkJCTxUeXBlPlRpbWU8L1R5cGU+CgkJCQk8UmFuZ2VFbnVtZXJhdGlvbiAvPgoJCQkJPFVuaXRzIC8+CgkJCQk8RGVzY3JpcHRpb24+PCFbQ0RBVEFbSW5kaWNhdGVzIHRoZSBEYXRhIGluc3RhbmNlIGNyZWF0aW9uIHRpbWVzdGFtcC5dXT48L0Rlc2NyaXB0aW9uPgoJCQk8L0l0ZW0+CgkJCTxJdGVtIElEPSIzIj48TmFtZT5EYXRhIERlc2NyaXB0aW9uPC9OYW1lPgoJCQkJPE9wZXJhdGlvbnM+Ulc8L09wZXJhdGlvbnM+CgkJCQk8TXVsdGlwbGVJbnN0YW5jZXM+U2luZ2xlPC9NdWx0aXBsZUluc3RhbmNlcz4KCQkJCTxNYW5kYXRvcnk+T3B0aW9uYWw8L01hbmRhdG9yeT4KCQkJCTxUeXBlPlN0cmluZzwvVHlwZT4KCQkJCTxSYW5nZUVudW1lcmF0aW9uPjMyIGJ5dGVzPC9SYW5nZUVudW1lcmF0aW9uPgoJCQkJPFVuaXRzIC8+CgkJCQk8RGVzY3JpcHRpb24+PCFbQ0RBVEFbSW5kaWNhdGVzIHRoZSBkYXRhIGRlc2NyaXB0aW9uLgplLmcuICJtZXRlciByZWFkaW5nIi5dXT48L0Rlc2NyaXB0aW9uPgoJCQk8L0l0ZW0+CgkJCTxJdGVtIElEPSI0Ij48TmFtZT5EYXRhIEZvcm1hdDwvTmFtZT4KCQkJCTxPcGVyYXRpb25zPlJXPC9PcGVyYXRpb25zPgoJCQkJPE11bHRpcGxlSW5zdGFuY2VzPlNpbmdsZTwvTXVsdGlwbGVJbnN0YW5jZXM+CgkJCQk8TWFuZGF0b3J5Pk9wdGlvbmFsPC9NYW5kYXRvcnk+CgkJCQk8VHlwZT5TdHJpbmc8L1R5cGU+CgkJCQk8UmFuZ2VFbnVtZXJhdGlvbj4zMiBieXRlczwvUmFuZ2VFbnVtZXJhdGlvbj4KCQkJCTxVbml0cyAvPgoJCQkJPERlc2NyaXB0aW9uPjwhW0NEQVRBW0luZGljYXRlcyB0aGUgZm9ybWF0IG9mIHRoZSBBcHBsaWNhdGlvbiBEYXRhLgplLmcuIFlHLU1ldGVyLVdhdGVyLVJlYWRpbmcKVVRGOC1zdHJpbmcKXV0+PC9EZXNjcmlwdGlvbj4KCQkJPC9JdGVtPgoJCQk8SXRlbSBJRD0iNSI+PE5hbWU+QXBwIElEPC9OYW1lPgoJCQkJPE9wZXJhdGlvbnM+Ulc8L09wZXJhdGlvbnM+CgkJCQk8TXVsdGlwbGVJbnN0YW5jZXM+U2luZ2xlPC9NdWx0aXBsZUluc3RhbmNlcz4KCQkJCTxNYW5kYXRvcnk+T3B0aW9uYWw8L01hbmRhdG9yeT4KCQkJCTxUeXBlPkludGVnZXI8L1R5cGU+CgkJCQk8UmFuZ2VFbnVtZXJhdGlvbj4yIGJ5dGVzPC9SYW5nZUVudW1lcmF0aW9uPgoJCQkJPFVuaXRzIC8+CgkJCQk8RGVzY3JpcHRpb24+PCFbQ0RBVEFbSW5kaWNhdGVzIHRoZSBkZXN0aW5hdGlvbiBBcHBsaWNhdGlvbiBJRC5dXT48L0Rlc2NyaXB0aW9uPgoJCQk8L0l0ZW0+PC9SZXNvdXJjZXM+CgkJPERlc2NyaXB0aW9uMj48IVtDREFUQVtdXT48L0Rlc2NyaXB0aW9uMj4KCTwvT2JqZWN0Pgo8L0xXTTJNPgo="; + private static final String FILE_NAME = "test.jks"; + + @Test + public void testResources_create_update_delete() throws Exception { + // create resource + TbResource resource = new TbResource(); + resource.setResourceType(ResourceType.JKS); + resource.setTitle("Edge Test Resource"); + resource.setFileName(FILE_NAME); + resource.setData(TEST_DATA); + + edgeImitator.expectMessageAmount(1); + TbResource savedResource = doPost("/api/resource", resource, TbResource.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof ResourceUpdateMsg); + ResourceUpdateMsg resourceUpdateMsg = (ResourceUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, resourceUpdateMsg.getMsgType()); + Assert.assertEquals(savedResource.getUuidId().getMostSignificantBits(), resourceUpdateMsg.getIdMSB()); + Assert.assertEquals(savedResource.getUuidId().getLeastSignificantBits(), resourceUpdateMsg.getIdLSB()); + Assert.assertEquals("Edge Test Resource", resourceUpdateMsg.getTitle()); + Assert.assertEquals(ResourceType.JKS.name(), resourceUpdateMsg.getResourceType()); + Assert.assertEquals(FILE_NAME, resourceUpdateMsg.getResourceKey()); + Assert.assertEquals(FILE_NAME, resourceUpdateMsg.getFileName()); + Assert.assertEquals(TEST_DATA, resourceUpdateMsg.getData()); + Assert.assertTrue(StringUtils.isNotBlank(resourceUpdateMsg.getEtag())); + + // update resource + edgeImitator.expectMessageAmount(1); + savedResource.setTitle("Updated Edge Test Resource"); + savedResource = doPost("/api/resource", savedResource, TbResource.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof ResourceUpdateMsg); + resourceUpdateMsg = (ResourceUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, resourceUpdateMsg.getMsgType()); + Assert.assertEquals("Updated Edge Test Resource", resourceUpdateMsg.getTitle()); + + // delete resource + edgeImitator.expectMessageAmount(1); + doDelete("/api/resource/" + savedResource.getUuidId()) + .andExpect(status().isOk()); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof ResourceUpdateMsg); + resourceUpdateMsg = (ResourceUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, resourceUpdateMsg.getMsgType()); + Assert.assertEquals(savedResource.getUuidId().getMostSignificantBits(), resourceUpdateMsg.getIdMSB()); + Assert.assertEquals(savedResource.getUuidId().getLeastSignificantBits(), resourceUpdateMsg.getIdLSB()); + } + @Test public void testSendResourceToCloud() throws Exception { UUID uuid = Uuids.timeBased(); @@ -41,9 +98,9 @@ public class ResourceEdgeTest extends AbstractEdgeTest { resourceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); resourceUpdateMsgBuilder.setTitle("Edge Test Resource"); resourceUpdateMsgBuilder.setResourceType(ResourceType.JKS.name()); - resourceUpdateMsgBuilder.setResourceKey("EdgeResource.jks"); - resourceUpdateMsgBuilder.setFileName("EdgeResource.jks"); - resourceUpdateMsgBuilder.setData("Data"); + resourceUpdateMsgBuilder.setResourceKey(FILE_NAME); + resourceUpdateMsgBuilder.setFileName(FILE_NAME); + resourceUpdateMsgBuilder.setData(TEST_DATA); resourceUpdateMsgBuilder.setIsSystem(false); resourceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE); testAutoGeneratedCodeByProtobuf(resourceUpdateMsgBuilder); diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index f9efbb3ee1..9be0627414 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -45,6 +45,7 @@ import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.v1.OtaPackageUpdateMsg; import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg; +import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.v1.TenantProfileUpdateMsg; @@ -303,6 +304,11 @@ public class EdgeImitator { result.add(saveDownlinkMsg(tenantProfileUpdateMsg)); } } + if (downlinkMsg.getResourceUpdateMsgCount() > 0) { + for (ResourceUpdateMsg resourceUpdateMsg : downlinkMsg.getResourceUpdateMsgList()) { + result.add(saveDownlinkMsg(resourceUpdateMsg)); + } + } if (downlinkMsg.hasEdgeConfiguration()) { result.add(saveDownlinkMsg(downlinkMsg.getEdgeConfiguration())); } From e6f33237bb80450b14038141e9342792b81aa613 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Tue, 31 Oct 2023 10:49:25 +0200 Subject: [PATCH 5/5] Refactor ResourceEdgeProcessor to use edgeId as sync instrument --- .../server/service/edge/rpc/EdgeGrpcSession.java | 2 +- .../edge/rpc/processor/resource/ResourceEdgeProcessor.java | 7 ++++--- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 81a77c1338..f9bd956706 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -715,7 +715,7 @@ public final class EdgeGrpcSession implements Closeable { } if (uplinkMsg.getResourceUpdateMsgCount() > 0) { for (ResourceUpdateMsg resourceUpdateMsg : uplinkMsg.getResourceUpdateMsgList()) { - result.add(ctx.getResourceEdgeProcessor().processResourceMsgFromEdge(edge.getTenantId(), resourceUpdateMsg)); + result.add(ctx.getResourceEdgeProcessor().processResourceMsgFromEdge(edge.getTenantId(), edge, resourceUpdateMsg)); } } if (uplinkMsg.getRuleChainMetadataRequestMsgCount() > 0) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java index ebc9a5440e..e7c5b21dbd 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/resource/ResourceEdgeProcessor.java @@ -21,6 +21,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.TbResourceId; import org.thingsboard.server.common.data.id.TenantId; @@ -37,10 +38,10 @@ import java.util.UUID; @TbCoreComponent public class ResourceEdgeProcessor extends BaseResourceProcessor { - public ListenableFuture processResourceMsgFromEdge(TenantId tenantId, ResourceUpdateMsg resourceUpdateMsg) { + public ListenableFuture processResourceMsgFromEdge(TenantId tenantId, Edge edge, ResourceUpdateMsg resourceUpdateMsg) { TbResourceId tbResourceId = new TbResourceId(new UUID(resourceUpdateMsg.getIdMSB(), resourceUpdateMsg.getIdLSB())); try { - edgeSynchronizationManager.getSync().set(true); + edgeSynchronizationManager.getEdgeId().set(edge.getId()); switch (resourceUpdateMsg.getMsgType()) { case ENTITY_CREATED_RPC_MESSAGE: @@ -64,7 +65,7 @@ public class ResourceEdgeProcessor extends BaseResourceProcessor { return Futures.immediateFailedFuture(e); } } finally { - edgeSynchronizationManager.getSync().remove(); + edgeSynchronizationManager.getEdgeId().remove(); } return Futures.immediateFuture(null); }