Browse Source

Added full or partial sync

pull/7395/head
Volodymyr Babak 4 years ago
parent
commit
7d90988058
  1. 4
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  3. 10
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  4. 32
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  5. 47
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EntityViewsEdgeEventFetcher.java
  6. 10
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java
  7. 2
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java
  8. 1
      common/edge-api/src/main/proto/edge.proto

4
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.device.DeviceService;
import org.thingsboard.server.dao.edge.EdgeEventService; import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.edge.EdgeService; 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.ota.OtaPackageService;
import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.queue.QueueService;
import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.rule.RuleChainService;
@ -88,6 +89,9 @@ public class EdgeContextComponent {
@Autowired @Autowired
private AssetService assetService; private AssetService assetService;
@Autowired
private EntityViewService entityViewService;
@Autowired @Autowired
private DeviceProfileService deviceProfileService; private DeviceProfileService deviceProfileService;

2
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) { if (session != null) {
boolean success = false; boolean success = false;
if (session.isConnected()) { if (session.isConnected()) {
session.startSyncProcess(tenantId, edgeId); session.startSyncProcess(tenantId, edgeId, true);
success = true; success = true;
} }
clusterService.pushEdgeSyncResponseToCore(new FromEdgeSyncResponse(requestId, tenantId, edgeId, success)); clusterService.pushEdgeSyncResponseToCore(new FromEdgeSyncResponse(requestId, tenantId, edgeId, success));

10
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 (connected) {
if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) {
if (requestMsg.hasSyncRequestMsg() && requestMsg.getSyncRequestMsg().getSyncRequired()) { 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 { } else {
syncCompleted = true; 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); log.trace("[{}][{}] Staring edge sync process", tenantId, edgeId);
syncCompleted = false; syncCompleted = false;
doSync(new EdgeSyncCursor(ctx, edge)); doSync(new EdgeSyncCursor(ctx, edge, fullSync));
} }
private void doSync(EdgeSyncCursor cursor) { private void doSync(EdgeSyncCursor cursor) {

32
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.DeviceProfilesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.DevicesEdgeEventFetcher; 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.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.OtaPackagesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.QueuesEdgeEventFetcher; 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.RuleChainsEdgeEventFetcher;
@ -44,23 +45,28 @@ public class EdgeSyncCursor {
int currentIdx = 0; int currentIdx = 0;
public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge) { public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge, boolean fullSync) {
fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService())); if (fullSync) {
fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService()));
fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), ctx.getFreemarkerConfig())); fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService()));
fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), ctx.getFreemarkerConfig()));
fetchers.add(new AssetProfilesEdgeEventFetcher(ctx.getAssetProfileService())); fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService()));
fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); fetchers.add(new AssetProfilesEdgeEventFetcher(ctx.getAssetProfileService()));
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService()));
fetchers.add(new CustomerEdgeEventFetcher()); if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) {
fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId())); fetchers.add(new CustomerEdgeEventFetcher());
fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId()));
}
} }
fetchers.add(new DevicesEdgeEventFetcher(ctx.getDeviceService())); fetchers.add(new DevicesEdgeEventFetcher(ctx.getDeviceService()));
fetchers.add(new AssetsEdgeEventFetcher(ctx.getAssetService())); fetchers.add(new AssetsEdgeEventFetcher(ctx.getAssetService()));
fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); fetchers.add(new EntityViewsEdgeEventFetcher(ctx.getEntityViewService()));
fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new DashboardsEdgeEventFetcher(ctx.getDashboardService())); 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() { public boolean hasNext() {

47
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<EntityView> {
private final EntityViewService entityViewService;
@Override
PageData<EntityView> 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);
}
}

10
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java

@ -207,9 +207,17 @@ public class EdgeGrpcClient implements EdgeRpcClient {
@Override @Override
public void sendSyncRequestMsg(boolean syncRequired) { public void sendSyncRequestMsg(boolean syncRequired) {
sendSyncRequestMsg(syncRequired, true);
}
@Override
public void sendSyncRequestMsg(boolean syncRequired, boolean fullSync) {
uplinkMsgLock.lock(); uplinkMsgLock.lock();
try { try {
SyncRequestMsg syncRequestMsg = SyncRequestMsg.newBuilder().setSyncRequired(syncRequired).build(); SyncRequestMsg syncRequestMsg = SyncRequestMsg.newBuilder()
.setSyncRequired(syncRequired)
.setFullSync(fullSync)
.build();
this.inputStream.onNext(RequestMsg.newBuilder() this.inputStream.onNext(RequestMsg.newBuilder()
.setMsgType(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE) .setMsgType(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)
.setSyncRequestMsg(syncRequestMsg) .setSyncRequestMsg(syncRequestMsg)

2
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);
void sendSyncRequestMsg(boolean syncRequired, boolean fullSync);
void sendUplinkMsg(UplinkMsg uplinkMsg); void sendUplinkMsg(UplinkMsg uplinkMsg);
void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg); void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg);

1
common/edge-api/src/main/proto/edge.proto

@ -84,6 +84,7 @@ message ConnectResponseMsg {
message SyncRequestMsg { message SyncRequestMsg {
bool syncRequired = 1; bool syncRequired = 1;
optional bool fullSync = 2;
} }
message SyncCompletedMsg { message SyncCompletedMsg {

Loading…
Cancel
Save