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 d0bba65837..de9d2d142a 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 @@ -48,6 +48,7 @@ import org.thingsboard.server.service.edge.rpc.processor.WidgetBundleEdgeProcess import org.thingsboard.server.service.edge.rpc.processor.WidgetTypeEdgeProcessor; import org.thingsboard.server.service.edge.rpc.sync.EdgeRequestsService; import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.executors.GrpcCallbackExecutorService; @Component @TbCoreComponent @@ -138,4 +139,7 @@ public class EdgeContextComponent { @Autowired private DbCallbackExecutorService dbCallbackExecutor; + + @Autowired + private GrpcCallbackExecutorService grpcCallbackExecutorService; } 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 cdb649037b..37d3c42c8a 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 @@ -16,8 +16,8 @@ package org.thingsboard.server.service.edge.rpc; import com.fasterxml.jackson.databind.ObjectMapper; -import com.google.common.io.Resources; import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; import io.grpc.Server; import io.grpc.netty.NettyServerBuilder; import io.grpc.stub.StreamObserver; @@ -46,7 +46,6 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.annotation.Nullable; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; -import java.io.File; import java.io.IOException; import java.io.InputStream; import java.util.Collections; @@ -93,6 +92,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Value("${edges.scheduler_pool_size}") private int schedulerPoolSize; + @Value("${edges.send_scheduler_pool_size}") + private int sendSchedulerPoolSize; + @Autowired private EdgeContextComponent ctx; @@ -105,6 +107,8 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private ExecutorService syncExecutorService; + private ScheduledExecutorService sendScheduler; + @PostConstruct public void init() { log.info("Initializing Edge RPC service!"); @@ -131,6 +135,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i throw new RuntimeException("Failed to start Edge RPC server!"); } this.scheduler = Executors.newScheduledThreadPool(schedulerPoolSize, ThingsBoardThreadFactory.forName("edge-scheduler")); + this.sendScheduler = Executors.newScheduledThreadPool(sendSchedulerPoolSize, ThingsBoardThreadFactory.forName("edge-send-scheduler")); this.syncExecutorService = Executors.newFixedThreadPool( Runtime.getRuntime().availableProcessors(), ThingsBoardThreadFactory.forName("edge-sync")); log.info("Edge RPC service initialized!"); @@ -152,6 +157,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (scheduler != null) { scheduler.shutdownNow(); } + if (sendScheduler != null) { + sendScheduler.shutdownNow(); + } if (syncExecutorService != null) { syncExecutorService.shutdownNow(); } @@ -159,7 +167,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public StreamObserver handleMsgs(StreamObserver outputStream) { - return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, mapper, syncExecutorService).getInputStream(); + return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, mapper, syncExecutorService, sendScheduler).getInputStream(); } @Override @@ -180,12 +188,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); session.close(); sessions.remove(edgeId); - Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); - lock.lock(); + final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + newEventLock.lock(); try { sessionNewEvents.remove(edgeId); } finally { - lock.unlock(); + newEventLock.unlock(); } cancelScheduleEdgeEventsCheck(edgeId); } @@ -194,27 +202,27 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void onEdgeEvent(TenantId tenantId, EdgeId edgeId) { log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId.getId()); - Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); - lock.lock(); + final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + newEventLock.lock(); try { if (Boolean.FALSE.equals(sessionNewEvents.get(edgeId))) { log.trace("[{}] set session new events flag to true [{}]", tenantId, edgeId.getId()); sessionNewEvents.put(edgeId, true); } } finally { - lock.unlock(); + newEventLock.unlock(); } } private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { log.info("[{}] edge [{}] connected successfully.", edgeGrpcSession.getSessionId(), edgeId); sessions.put(edgeId, edgeGrpcSession); - Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); - lock.lock(); + final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + newEventLock.lock(); try { sessionNewEvents.put(edgeId, true); } finally { - lock.unlock(); + newEventLock.unlock(); } save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true); save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis()); @@ -239,21 +247,33 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (sessions.containsKey(edgeId)) { ScheduledFuture schedule = scheduler.schedule(() -> { try { - Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); - lock.lock(); + final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + newEventLock.lock(); try { if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { log.trace("[{}] Set session new events flag to false", edgeId.getId()); sessionNewEvents.put(edgeId, false); - session.processEdgeEvents(); + Futures.addCallback(session.processEdgeEvents(), new FutureCallback<>() { + @Override + public void onSuccess(Void result) { + scheduleEdgeEventsCheck(session); + } + + @Override + public void onFailure(Throwable t) { + log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), t); + scheduleEdgeEventsCheck(session); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + scheduleEdgeEventsCheck(session); } } finally { - lock.unlock(); + newEventLock.unlock(); } } catch (Exception e) { log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), e); } - scheduleEdgeEventsCheck(session); }, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS); sessionEdgeEventChecks.put(edgeId, schedule); log.trace("[{}] Check edge event scheduled for edge [{}]", tenantId, edgeId.getId()); @@ -277,12 +297,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void onEdgeDisconnect(EdgeId edgeId) { log.info("[{}] edge disconnected!", edgeId); sessions.remove(edgeId); - Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); - lock.lock(); + final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + newEventLock.lock(); try { sessionNewEvents.remove(edgeId); } finally { - lock.unlock(); + newEventLock.unlock(); } save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis()); 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 aca42db0b4..f16df05af8 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 @@ -20,7 +20,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; +import com.google.common.util.concurrent.SettableFuture; import io.grpc.stub.StreamObserver; import lombok.Data; import lombok.extern.slf4j.Slf4j; @@ -31,9 +31,7 @@ import org.thingsboard.server.common.data.EdgeUtils; 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.EdgeId; -import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; @@ -69,32 +67,21 @@ import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; -import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.AssetsEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.CustomerUsersEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.GeneralEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetsBundlesEdgeEventFetcher; import java.io.Closeable; import java.util.ArrayList; -import java.util.Collection; import java.util.Collections; -import java.util.HashMap; -import java.util.Iterator; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.UUID; -import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantLock; import java.util.function.BiConsumer; @@ -116,6 +103,8 @@ public final class EdgeGrpcSession implements Closeable { private final ObjectMapper mapper; private final Map pendingMsgsMap = new LinkedHashMap<>(); + private SettableFuture pendingFuture; + private ScheduledFuture sendSchedule; private EdgeContextComponent ctx; private Edge edge; @@ -126,10 +115,10 @@ public final class EdgeGrpcSession implements Closeable { private ExecutorService syncExecutorService; - private CountDownLatch latch; + private ScheduledExecutorService sendScheduler; EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, BiConsumer sessionOpenListener, - Consumer sessionCloseListener, ObjectMapper mapper, ExecutorService syncExecutorService) { + Consumer sessionCloseListener, ObjectMapper mapper, ExecutorService syncExecutorService, ScheduledExecutorService sendScheduler) { this.sessionId = UUID.randomUUID(); this.ctx = ctx; this.outputStream = outputStream; @@ -137,6 +126,7 @@ public final class EdgeGrpcSession implements Closeable { this.sessionCloseListener = sessionCloseListener; this.mapper = mapper; this.syncExecutorService = syncExecutorService; + this.sendScheduler = sendScheduler; initInputStream(); } @@ -155,19 +145,21 @@ public final class EdgeGrpcSession implements Closeable { connected = true; } } - if (connected && requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { - if (requestMsg.getSyncRequestMsg().getSyncRequired()) { - startSyncProcess(edge.getTenantId(), edge.getId()); - } else { - syncCompleted = true; - } - } if (connected) { - if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE) && requestMsg.hasUplinkMsg()) { - onUplinkMsg(requestMsg.getUplinkMsg()); + if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { + if (requestMsg.hasSyncRequestMsg() && requestMsg.getSyncRequestMsg().getSyncRequired()) { + startSyncProcess(edge.getTenantId(), edge.getId()); + } else { + syncCompleted = true; + } } - if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE) && requestMsg.hasDownlinkResponseMsg()) { - onDownlinkResponse(requestMsg.getDownlinkResponseMsg()); + if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) { + if (requestMsg.hasUplinkMsg()) { + onUplinkMsg(requestMsg.getUplinkMsg()); + } + if (requestMsg.hasDownlinkResponseMsg()) { + onDownlinkResponse(requestMsg.getDownlinkResponseMsg()); + } } } } @@ -204,39 +196,47 @@ public final class EdgeGrpcSession implements Closeable { syncCompleted = false; syncExecutorService.submit(() -> { try { - startProcessingEdgeEvents(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); - startProcessingEdgeEvents(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); - startProcessingEdgeEvents(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); - startProcessingEdgeEvents(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); - - startProcessingEdgeEvents(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); - if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { - EdgeEvent customerEdgeEvent = EdgeEventUtils.constructEdgeEvent(edge.getTenantId(), edge.getId(), - EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null); - DownlinkMsg customerDownlinkMsg = convertToDownlinkMsg(customerEdgeEvent); - sendDownlinkMsgsPack(Collections.singletonList(customerDownlinkMsg)); - - startProcessingEdgeEvents(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId())); - } - - startProcessingEdgeEvents(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService())); - - startProcessingEdgeEvents(new AssetsEdgeEventFetcher(ctx.getAssetService())); - startProcessingEdgeEvents(new DashboardsEdgeEventFetcher(ctx.getDashboardService())); - - DownlinkMsg syncCompleteDownlinkMsg = DownlinkMsg.newBuilder() - .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .setSyncCompletedMsg(SyncCompletedMsg.newBuilder().build()) - .build(); - sendDownlinkMsgsPack(Collections.singletonList(syncCompleteDownlinkMsg)); - - syncCompleted = true; + doSync(new EdgeSyncCursor(ctx, edge)); } catch (Exception e) { log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), e); } }); } + private void doSync(EdgeSyncCursor cursor) { + if (cursor.hasNext()) { + log.info("[{}][{}] starting sync process, cursor current idx = {}", edge.getTenantId(), edge.getId(), cursor.getCurrentIdx()); + ListenableFuture uuidListenableFuture = startProcessingEdgeEvents(cursor.getNext()); + Futures.addCallback(uuidListenableFuture, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable UUID result) { + doSync(cursor); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), t); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + DownlinkMsg syncCompleteDownlinkMsg = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .setSyncCompletedMsg(SyncCompletedMsg.newBuilder().build()) + .build(); + Futures.addCallback(sendDownlinkMsgsPack(Collections.singletonList(syncCompleteDownlinkMsg)), new FutureCallback() { + @Override + public void onSuccess(Void result) { + syncCompleted = true; + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Exception during sending sync complete", edge.getTenantId(), edge.getId(), t); + } + }, ctx.getGrpcCallbackExecutorService()); + } + } + private void onUplinkMsg(UplinkMsg uplinkMsg) { ListenableFuture> future = processUplinkMsg(uplinkMsg); Futures.addCallback(future, new FutureCallback<>() { @@ -260,7 +260,7 @@ public final class EdgeGrpcSession implements Closeable { .setUplinkResponseMsg(uplinkResponseMsg) .build()); } - }, MoreExecutors.directExecutor()); + }, ctx.getGrpcCallbackExecutorService()); } private void onDownlinkResponse(DownlinkResponseMsg msg) { @@ -271,7 +271,13 @@ public final class EdgeGrpcSession implements Closeable { } else { log.error("[{}] Msg processing failed! Error msg: {}", edge.getRoutingKey(), msg.getErrorMsg()); } - latch.countDown(); + if (pendingMsgsMap.size() == 0) { + log.debug("[{}] Pending msgs map is empty. Stopping current iteration {}", edge.getRoutingKey(), msg); + if (sendSchedule != null) { + sendSchedule.cancel(false); + } + pendingFuture.set(null); + } } catch (Exception e) { log.error("[{}] Can't process downlink response message [{}]", this.sessionId, msg, e); } @@ -305,76 +311,136 @@ public final class EdgeGrpcSession implements Closeable { sendDownlinkMsg(edgeConfigMsg); } - void processEdgeEvents() throws Exception { - log.trace("[{}] processHandleMessages started", this.sessionId); + ListenableFuture processEdgeEvents() throws Exception { + SettableFuture result = SettableFuture.create(); + log.trace("[{}] starting processing edge events", this.sessionId); if (isConnected() && isSyncCompleted()) { Long queueStartTs = getQueueStartTs().get(); GeneralEdgeEventFetcher fetcher = new GeneralEdgeEventFetcher( queueStartTs, ctx.getEdgeEventService()); - UUID ifOffset = startProcessingEdgeEvents(fetcher); - if (ifOffset != null) { - Long newStartTs = Uuids.unixTimestamp(ifOffset); - updateQueueStartTs(newStartTs); - log.debug("[{}] queue offset was updated [{}][{}]", this.sessionId, ifOffset, newStartTs); - } + ListenableFuture ifOffsetFuture = startProcessingEdgeEvents(fetcher); + Futures.addCallback(ifOffsetFuture, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable UUID ifOffset) { + if (ifOffset != null) { + Long newStartTs = Uuids.unixTimestamp(ifOffset); + ListenableFuture> updateFuture = updateQueueStartTs(newStartTs); + Futures.addCallback(updateFuture, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable List list) { + log.debug("[{}] queue offset was updated [{}][{}]", sessionId, ifOffset, newStartTs); + result.set(null); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}] Failed to update queue offset [{}]", sessionId, ifOffset, t); + result.setException(t); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + log.trace("[{}] ifOffset is null. Skipping iteration without db update", sessionId); + result.set(null); + } + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}] Failed to process events", sessionId, t); + result.setException(t); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + log.trace("[{}] edge is not connected or sync is not completed. Skipping iteration", sessionId); + result.set(null); } - log.trace("[{}] processHandleMessages finished", this.sessionId); + return result; } - private UUID startProcessingEdgeEvents(EdgeEventFetcher fetcher) throws Exception { + private ListenableFuture startProcessingEdgeEvents(EdgeEventFetcher fetcher) { + SettableFuture result = SettableFuture.create(); PageLink pageLink = fetcher.getPageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()); - PageData pageData; - UUID ifOffset = null; - boolean success; - do { - pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); + processEdgeEvents(fetcher, pageLink, result); + return result; + } + + private void processEdgeEvents(EdgeEventFetcher fetcher, PageLink pageLink, SettableFuture result) { + try { + PageData pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); if (isConnected() && !pageData.getData().isEmpty()) { log.trace("[{}] [{}] event(s) are going to be processed.", this.sessionId, pageData.getData().size()); List downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); - success = sendDownlinkMsgsPack(downlinkMsgsPack); - ifOffset = pageData.getData().get(pageData.getData().size() - 1).getUuidId(); - if (success) { - pageLink = pageLink.nextPageLink(); - } + Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback() { + @Override + public void onSuccess(@Nullable Void tmp) { + if (isConnected() && pageData.hasNext()) { + processEdgeEvents(fetcher, pageLink.nextPageLink(), result); + } else { + UUID ifOffset = pageData.getData().get(pageData.getData().size() - 1).getUuidId(); + result.set(ifOffset); + } + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}] Failed to send downlink msgs pack", sessionId, t); + result.setException(t); + } + }, ctx.getGrpcCallbackExecutorService()); } else { log.trace("[{}] no event(s) found. Stop processing edge events", this.sessionId); + result.set(null); } - } while (isConnected() && pageData.hasNext()); - return ifOffset; + } catch (Exception e) { + log.error("[{}] Failed to fetch edge events", this.sessionId, e); + result.setException(e); + } } - private boolean sendDownlinkMsgsPack(List downlinkMsgsPack) throws InterruptedException { + private ListenableFuture sendDownlinkMsgsPack(List downlinkMsgsPack) { + SettableFuture result = SettableFuture.create(); downlinkMsgsPackLock.lock(); try { - boolean success; pendingMsgsMap.clear(); downlinkMsgsPack.forEach(msg -> pendingMsgsMap.put(msg.getDownlinkMsgId(), msg)); - do { - log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, pendingMsgsMap.values().size()); - latch = new CountDownLatch(pendingMsgsMap.values().size()); - List copy = new ArrayList<>(pendingMsgsMap.values()); - for (DownlinkMsg downlinkMsg : copy) { - sendDownlinkMsg(ResponseMsg.newBuilder() - .setDownlinkMsg(downlinkMsg) - .build()); - } - success = latch.await(10, TimeUnit.SECONDS); - if (!success || pendingMsgsMap.values().size() > 0) { - log.warn("[{}] Failed to deliver the batch: {}", this.sessionId, pendingMsgsMap.values()); - } - if (isConnected() && (!success || pendingMsgsMap.values().size() > 0)) { - try { - Thread.sleep(ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches()); - } catch (InterruptedException e) { - log.error("[{}] Error during sleep between batches", this.sessionId, e); - } - } - } while (isConnected() && (!success || pendingMsgsMap.values().size() > 0)); - return success; + scheduleDownlinkMsgsPackSend(true, result); } finally { downlinkMsgsPackLock.unlock(); } + return result; + } + + private void scheduleDownlinkMsgsPackSend(boolean firstRun, SettableFuture futureResult) { + pendingFuture = futureResult; + Runnable runnable = () -> { + try { + if (isConnected() && pendingMsgsMap.values().size() > 0) { + if (!firstRun) { + log.warn("[{}] Failed to deliver the batch: {}", this.sessionId, pendingMsgsMap.values()); + } + log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, pendingMsgsMap.values().size()); + List copy = new ArrayList<>(pendingMsgsMap.values()); + for (DownlinkMsg downlinkMsg : copy) { + sendDownlinkMsg(ResponseMsg.newBuilder() + .setDownlinkMsg(downlinkMsg) + .build()); + } + scheduleDownlinkMsgsPackSend(false, futureResult); + } else { + futureResult.set(null); + } + } catch (Exception e) { + futureResult.setException(e); + } + }; + + if (firstRun) { + syncExecutorService.submit(runnable); + } else { + sendSchedule = sendScheduler.schedule(runnable, ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches(), TimeUnit.MILLISECONDS); + } + } private DownlinkMsg convertToDownlinkMsg(EdgeEvent edgeEvent) { @@ -439,14 +505,14 @@ public final class EdgeGrpcSession implements Closeable { } else { return 0L; } - }, ctx.getDbCallbackExecutor()); + }, ctx.getGrpcCallbackExecutorService()); } - private void updateQueueStartTs(Long newStartTs) { + private ListenableFuture> updateQueueStartTs(Long newStartTs) { log.trace("[{}] updating QueueStartTs [{}][{}]", this.sessionId, edge.getId(), newStartTs); newStartTs = ++newStartTs; // increments ts by 1 - next edge event search starts from current offset + 1 List attributes = Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_TS_ATTR_KEY, newStartTs), System.currentTimeMillis())); - ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); + return ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); } private DownlinkMsg processEntityMessage(EdgeEvent edgeEvent, EdgeEventActionType action) { 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 new file mode 100644 index 0000000000..85ef95ee12 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java @@ -0,0 +1,74 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.edge.rpc; + +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.service.edge.EdgeContextComponent; +import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.AssetsEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.CustomerEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.CustomerUsersEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetsBundlesEdgeEventFetcher; + +import java.util.LinkedList; +import java.util.List; + +public class EdgeSyncCursor { + + List fetchers = new LinkedList<>(); + + int currentIdx = 0; + + public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge) { + fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); + fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); + fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); + fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); + fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); + if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { + fetchers.add(new CustomerEdgeEventFetcher()); + fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId())); + } + fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService())); + fetchers.add(new AssetsEdgeEventFetcher(ctx.getAssetService())); + fetchers.add(new DashboardsEdgeEventFetcher(ctx.getDashboardService())); + } + + public boolean hasNext() { + return fetchers.size() > currentIdx; + } + + public EdgeEventFetcher getNext() { + if (hasNext()) { + EdgeEventFetcher edgeEventFetcher = fetchers.get(currentIdx); + currentIdx++; + return edgeEventFetcher; + } else { + throw new IndexOutOfBoundsException(); + } + } + + public int getCurrentIdx() { + return currentIdx; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java index c59d6c5001..0035fd1334 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java @@ -56,7 +56,7 @@ public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher { @Override public PageLink getPageLink(int pageSize) { - return new PageLink(pageSize); + return null; } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerEdgeEventFetcher.java new file mode 100644 index 0000000000..abdde978f4 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerEdgeEventFetcher.java @@ -0,0 +1,49 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.edge.rpc.fetch; + +import lombok.AllArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.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.service.edge.rpc.EdgeEventUtils; + +import java.util.ArrayList; +import java.util.List; + +@AllArgsConstructor +@Slf4j +public class CustomerEdgeEventFetcher implements EdgeEventFetcher { + + @Override + public PageLink getPageLink(int pageSize) { + return null; + } + + @Override + public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { + List result = new ArrayList<>(); + result.add(EdgeEventUtils.constructEdgeEvent(edge.getTenantId(), edge.getId(), + EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null)); + // @voba - returns PageData object to be in sync with other fetchers + return new PageData<>(result, 1, result.size(), false); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/executors/GrpcCallbackExecutorService.java b/application/src/main/java/org/thingsboard/server/service/executors/GrpcCallbackExecutorService.java new file mode 100644 index 0000000000..2de1c21974 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/executors/GrpcCallbackExecutorService.java @@ -0,0 +1,33 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.executors; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.common.util.AbstractListeningExecutor; + +@Component +public class GrpcCallbackExecutorService extends AbstractListeningExecutor { + + @Value("${edges.grpc_callback_thread_pool_size}") + private int grpcCallbackExecutorThreadPoolSize; + + @Override + protected int getThreadPollSize() { + return grpcCallbackExecutorThreadPoolSize; + } + +} diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 9444002bf7..89adef5456 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -707,7 +707,9 @@ edges: max_read_records_count: "${EDGES_STORAGE_MAX_READ_RECORDS_COUNT:50}" no_read_records_sleep: "${EDGES_NO_READ_RECORDS_SLEEP:1000}" sleep_between_batches: "${EDGES_SLEEP_BETWEEN_BATCHES:1000}" - scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:4}" + scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}" + send_scheduler_pool_size: "${EDGES_SEND_SCHEDULER_POOL_SIZE:1}" + grpc_callback_thread_pool_size: "${EDGES_GRPC_CALLBACK_POOL_SIZE:1}" edge_events_ttl: "${EDGES_EDGE_EVENTS_TTL:0}" state: persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}"