From 7d909880589977b772d658f19c60bd256dc14524 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 19 Oct 2022 12:42:09 +0300 Subject: [PATCH] Added full or partial sync --- .../service/edge/EdgeContextComponent.java | 4 ++ .../service/edge/rpc/EdgeGrpcService.java | 2 +- .../service/edge/rpc/EdgeGrpcSession.java | 10 ++-- .../service/edge/rpc/EdgeSyncCursor.java | 32 ++++++++----- .../fetch/EntityViewsEdgeEventFetcher.java | 47 +++++++++++++++++++ .../thingsboard/edge/rpc/EdgeGrpcClient.java | 10 +++- .../thingsboard/edge/rpc/EdgeRpcClient.java | 2 + common/edge-api/src/main/proto/edge.proto | 1 + 8 files changed, 90 insertions(+), 18 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EntityViewsEdgeEventFetcher.java 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 5a0d52d734..4adb73ebbb 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 @@ -30,6 +30,7 @@ import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.edge.EdgeEventService; 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.rule.RuleChainService; @@ -88,6 +89,9 @@ public class EdgeContextComponent { @Autowired private AssetService assetService; + @Autowired + private EntityViewService entityViewService; + @Autowired private DeviceProfileService deviceProfileService; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index 4cb6fcdb0a..fea89d6336 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java @@ -271,7 +271,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (session != null) { boolean success = false; if (session.isConnected()) { - session.startSyncProcess(tenantId, edgeId); + session.startSyncProcess(tenantId, edgeId, true); success = true; } clusterService.pushEdgeSyncResponseToCore(new FromEdgeSyncResponse(requestId, tenantId, edgeId, success)); 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 562d34012c..9864d89048 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 @@ -134,7 +134,11 @@ public final class EdgeGrpcSession implements Closeable { if (connected) { if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { if (requestMsg.hasSyncRequestMsg() && requestMsg.getSyncRequestMsg().getSyncRequired()) { - startSyncProcess(edge.getTenantId(), edge.getId()); + boolean fullSync = true; + if (requestMsg.getSyncRequestMsg().hasFullSync()) { + fullSync = requestMsg.getSyncRequestMsg().getFullSync(); + } + startSyncProcess(edge.getTenantId(), edge.getId(), fullSync); } else { syncCompleted = true; } @@ -177,10 +181,10 @@ public final class EdgeGrpcSession implements Closeable { }; } - public void startSyncProcess(TenantId tenantId, EdgeId edgeId) { + public void startSyncProcess(TenantId tenantId, EdgeId edgeId, boolean fullSync) { log.trace("[{}][{}] Staring edge sync process", tenantId, edgeId); syncCompleted = false; - doSync(new EdgeSyncCursor(ctx, edge)); + doSync(new EdgeSyncCursor(ctx, edge, fullSync)); } private void doSync(EdgeSyncCursor cursor) { 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 b52370adec..e9232584be 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 @@ -27,6 +27,7 @@ import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.DevicesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher; +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; @@ -44,23 +45,28 @@ public class EdgeSyncCursor { int currentIdx = 0; - public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge) { - fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService())); - fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); - fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), ctx.getFreemarkerConfig())); - fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); - fetchers.add(new AssetProfilesEdgeEventFetcher(ctx.getAssetProfileService())); - fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); - if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { - fetchers.add(new CustomerEdgeEventFetcher()); - fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId())); + public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge, boolean fullSync) { + if (fullSync) { + fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService())); + fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); + fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), ctx.getFreemarkerConfig())); + fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); + fetchers.add(new AssetProfilesEdgeEventFetcher(ctx.getAssetProfileService())); + fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); + if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { + fetchers.add(new CustomerEdgeEventFetcher()); + fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId())); + } } fetchers.add(new DevicesEdgeEventFetcher(ctx.getDeviceService())); fetchers.add(new AssetsEdgeEventFetcher(ctx.getAssetService())); - fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); - fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); + fetchers.add(new EntityViewsEdgeEventFetcher(ctx.getEntityViewService())); fetchers.add(new DashboardsEdgeEventFetcher(ctx.getDashboardService())); - fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService())); + if (fullSync) { + fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); + fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); + fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService())); + } } public boolean hasNext() { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EntityViewsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EntityViewsEdgeEventFetcher.java new file mode 100644 index 0000000000..3bc4befdbe --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EntityViewsEdgeEventFetcher.java @@ -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.EntityView; +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.entityview.EntityViewService; + +@AllArgsConstructor +@Slf4j +public class EntityViewsEdgeEventFetcher extends BasePageableEdgeEventFetcher { + + private final EntityViewService entityViewService; + + @Override + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return entityViewService.findEntityViewsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, EntityView entityView) { + return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ENTITY_VIEW, + EdgeEventActionType.ADDED, entityView.getId(), null); + } +} diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java index 119b011adb..b6d9464bdd 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java @@ -207,9 +207,17 @@ public class EdgeGrpcClient implements EdgeRpcClient { @Override public void sendSyncRequestMsg(boolean syncRequired) { + sendSyncRequestMsg(syncRequired, true); + } + + @Override + public void sendSyncRequestMsg(boolean syncRequired, boolean fullSync) { uplinkMsgLock.lock(); try { - SyncRequestMsg syncRequestMsg = SyncRequestMsg.newBuilder().setSyncRequired(syncRequired).build(); + SyncRequestMsg syncRequestMsg = SyncRequestMsg.newBuilder() + .setSyncRequired(syncRequired) + .setFullSync(fullSync) + .build(); this.inputStream.onNext(RequestMsg.newBuilder() .setMsgType(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE) .setSyncRequestMsg(syncRequestMsg) diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java index 835f865e05..08214c61bd 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java @@ -36,6 +36,8 @@ public interface EdgeRpcClient { void sendSyncRequestMsg(boolean syncRequired); + void sendSyncRequestMsg(boolean syncRequired, boolean fullSync); + void sendUplinkMsg(UplinkMsg uplinkMsg); void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg); diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 6333184c56..a2e3b5989b 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -84,6 +84,7 @@ message ConnectResponseMsg { message SyncRequestMsg { bool syncRequired = 1; + optional bool fullSync = 2; } message SyncCompletedMsg {