diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java deleted file mode 100644 index 9fd94b4988..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java +++ /dev/null @@ -1,939 +0,0 @@ -/** - * Copyright © 2016-2024 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 com.datastax.oss.driver.api.core.uuid.Uuids; -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.SettableFuture; -import io.grpc.stub.StreamObserver; -import lombok.Data; -import lombok.extern.slf4j.Slf4j; -import org.checkerframework.checker.nullness.qual.Nullable; -import org.springframework.data.util.Pair; -import org.thingsboard.server.common.data.AttributeScope; -import org.thingsboard.server.common.data.DataConstants; -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.id.EdgeId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.AttributeKvEntry; -import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; -import org.thingsboard.server.common.data.kv.LongDataEntry; -import org.thingsboard.server.common.data.kv.StringDataEntry; -import org.thingsboard.server.common.data.limit.LimitedApi; -import org.thingsboard.server.common.data.notification.rule.trigger.EdgeCommunicationFailureTrigger; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.page.SortOrder; -import org.thingsboard.server.common.data.page.TimePageLink; -import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg; -import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; -import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; -import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; -import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; -import org.thingsboard.server.gen.edge.v1.AttributesRequestMsg; -import org.thingsboard.server.gen.edge.v1.ConnectRequestMsg; -import org.thingsboard.server.gen.edge.v1.ConnectResponseCode; -import org.thingsboard.server.gen.edge.v1.ConnectResponseMsg; -import org.thingsboard.server.gen.edge.v1.DashboardUpdateMsg; -import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; -import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; -import org.thingsboard.server.gen.edge.v1.DeviceProfileUpdateMsg; -import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg; -import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; -import org.thingsboard.server.gen.edge.v1.DownlinkMsg; -import org.thingsboard.server.gen.edge.v1.DownlinkResponseMsg; -import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; -import org.thingsboard.server.gen.edge.v1.EdgeUpdateMsg; -import org.thingsboard.server.gen.edge.v1.EdgeVersion; -import org.thingsboard.server.gen.edge.v1.EntityDataProto; -import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; -import org.thingsboard.server.gen.edge.v1.EntityViewsRequestMsg; -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; -import org.thingsboard.server.gen.edge.v1.UplinkMsg; -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.EdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.fetch.GeneralEdgeEventFetcher; -import org.thingsboard.server.service.edge.rpc.processor.alarm.AlarmProcessor; -import org.thingsboard.server.service.edge.rpc.processor.asset.AssetProcessor; -import org.thingsboard.server.service.edge.rpc.processor.asset.profile.AssetProfileProcessor; -import org.thingsboard.server.service.edge.rpc.processor.dashboard.DashboardProcessor; -import org.thingsboard.server.service.edge.rpc.processor.device.DeviceProcessor; -import org.thingsboard.server.service.edge.rpc.processor.device.profile.DeviceProfileProcessor; -import org.thingsboard.server.service.edge.rpc.processor.entityview.EntityViewProcessor; -import org.thingsboard.server.service.edge.rpc.processor.relation.RelationProcessor; -import org.thingsboard.server.service.edge.rpc.processor.resource.ResourceProcessor; - -import java.io.Closeable; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.Collections; -import java.util.List; -import java.util.Optional; -import java.util.UUID; -import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.locks.ReentrantLock; -import java.util.function.BiConsumer; - -@Slf4j -@Data -public abstract class AbstractEdgeGrpcSession implements EdgeGrpcSession, Closeable { - - private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; - private static final String QUEUE_START_SEQ_ID_ATTR_KEY = "queueStartSeqId"; - - private static final int MAX_DOWNLINK_ATTEMPTS = 10; - private static final String RATE_LIMIT_REACHED = "Rate limit reached"; - - protected static final ConcurrentLinkedQueue highPriorityQueue = new ConcurrentLinkedQueue<>(); - - protected UUID sessionId; - private BiConsumer sessionOpenListener; - private BiConsumer sessionCloseListener; - - private final EdgeSessionState sessionState = new EdgeSessionState(); - private final ReentrantLock downlinkMsgLock = new ReentrantLock(); - - protected EdgeContextComponent ctx; - protected Edge edge; - protected TenantId tenantId; - - private Long newStartTs; - private Long previousStartTs; - private Long newStartSeqId; - private Long previousStartSeqId; - private Long seqIdEnd; - - private StreamObserver inputStream; - private StreamObserver outputStream; - - private boolean connected; - private volatile boolean syncCompleted; - - private EdgeVersion edgeVersion; - private int maxInboundMessageSize; - private int clientMaxInboundMessageSize; - private int maxHighPriorityQueueSizePerSession; - - private ScheduledExecutorService sendDownlinkExecutorService; - - public AbstractEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, - BiConsumer sessionOpenListener, - BiConsumer sessionCloseListener, - ScheduledExecutorService sendDownlinkExecutorService, - int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { - this.sessionId = UUID.randomUUID(); - this.ctx = ctx; - this.outputStream = outputStream; - this.sessionOpenListener = sessionOpenListener; - this.sessionCloseListener = sessionCloseListener; - this.sendDownlinkExecutorService = sendDownlinkExecutorService; - this.maxInboundMessageSize = maxInboundMessageSize; - this.maxHighPriorityQueueSizePerSession = maxHighPriorityQueueSizePerSession; - initInputStream(); - } - - public void initInputStream() { - inputStream = new StreamObserver<>() { - @Override - public void onNext(RequestMsg requestMsg) { - if (!connected && requestMsg.getMsgType().equals(RequestMsgType.CONNECT_RPC_MESSAGE)) { - ConnectResponseMsg responseMsg = processConnect(requestMsg.getConnectRequestMsg()); - outputStream.onNext(ResponseMsg.newBuilder() - .setConnectResponseMsg(responseMsg) - .build()); - if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) { - outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); - } else { - if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) { - log.debug("[{}][{}] Client max inbound message size: {}", tenantId, sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); - clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); - } - connected = true; - } - } - if (connected) { - if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { - if (requestMsg.hasSyncRequestMsg()) { - boolean fullSync = false; - if (requestMsg.getSyncRequestMsg().hasFullSync()) { - fullSync = requestMsg.getSyncRequestMsg().getFullSync(); - } - startSyncProcess(fullSync); - } else { - syncCompleted = true; - } - } - if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) { - if (requestMsg.hasUplinkMsg()) { - onUplinkMsg(requestMsg.getUplinkMsg()); - } - if (requestMsg.hasDownlinkResponseMsg()) { - onDownlinkResponse(requestMsg.getDownlinkResponseMsg()); - } - } - } - } - - @Override - public void onError(Throwable t) { - log.error("[{}][{}] Stream was terminated due to error:", tenantId, sessionId, t); - closeSession(); - } - - @Override - public void onCompleted() { - log.info("[{}][{}] Stream was closed and completed successfully!", tenantId, sessionId); - closeSession(); - } - - private void closeSession() { - connected = false; - if (edge != null) { - try { - sessionCloseListener.accept(edge, sessionId); - } catch (Exception ignored) { - } - } - try { - outputStream.onCompleted(); - } catch (Exception ignored) { - } - } - }; - } - - @Override - public void onConfigurationUpdate(Edge edge) { - log.debug("[{}] onConfigurationUpdate [{}]", sessionId, edge); - this.tenantId = edge.getTenantId(); - this.edge = edge; - EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() - .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)).build(); - ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder() - .setEdgeUpdateMsg(edgeConfig) - .build(); - sendDownlinkMsg(edgeConfigMsg); - } - - @Override - public void startSyncProcess(boolean fullSync) { - log.info("[{}][{}][{}] Staring edge sync process", tenantId, edge.getId(), sessionId); - syncCompleted = false; - interruptGeneralProcessingOnSync(); - doSync(new EdgeSyncCursor(ctx, edge, fullSync)); - } - - private void doSync(EdgeSyncCursor cursor) { - if (cursor.hasNext()) { - EdgeEventFetcher next = cursor.getNext(); - log.info("[{}][{}] starting sync process, cursor current idx = {}, class = {}", - tenantId, edge.getId(), cursor.getCurrentIdx(), next.getClass().getSimpleName()); - ListenableFuture> future = startProcessingEdgeEvents(next); - Futures.addCallback(future, new FutureCallback<>() { - @Override - public void onSuccess(@Nullable Pair result) { - doSync(cursor); - } - - @Override - public void onFailure(Throwable t) { - log.error("[{}][{}] Exception during sync process", tenantId, edge.getId(), t); - } - }, ctx.getGrpcCallbackExecutorService()); - } else { - log.info("[{}][{}] sync process completed", tenantId, edge.getId()); - DownlinkMsg syncCompleteDownlinkMsg = DownlinkMsg.newBuilder() - .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .setSyncCompletedMsg(SyncCompletedMsg.newBuilder().build()) - .build(); - Futures.addCallback(sendDownlinkMsgsPack(Collections.singletonList(syncCompleteDownlinkMsg)), new FutureCallback<>() { - @Override - public void onSuccess(Boolean isInterrupted) { - markSyncCompletedSendEdgeEventUpdate(); - } - - @Override - public void onFailure(Throwable t) { - log.error("[{}][{}] Exception during sending sync complete", tenantId, edge.getId(), t); - markSyncCompletedSendEdgeEventUpdate(); - } - }, ctx.getGrpcCallbackExecutorService()); - } - } - - protected void processEdgeEvents(EdgeEventFetcher fetcher, PageLink pageLink, SettableFuture> result) { - try { - if (!highPriorityQueue.isEmpty()) { - processHighPriorityEvents(); - } - PageData pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); - if (isConnected() && !pageData.getData().isEmpty()) { - log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, sessionId, pageData.getData().size()); - List downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); - Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { - @Override - public void onSuccess(@Nullable Boolean isInterrupted) { - if (Boolean.TRUE.equals(isInterrupted)) { - log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId); - result.set(null); - } else { - if (isConnected() && pageData.hasNext()) { - processEdgeEvents(fetcher, pageLink.nextPageLink(), result); - } else { - EdgeEvent latestEdgeEvent = pageData.getData().get(pageData.getData().size() - 1); - UUID idOffset = latestEdgeEvent.getUuidId(); - if (idOffset != null) { - Long newStartTs = Uuids.unixTimestamp(idOffset); - long newStartSeqId = latestEdgeEvent.getSeqId(); - result.set(Pair.of(newStartTs, newStartSeqId)); - } else { - result.set(null); - } - } - } - } - - @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", sessionId); - result.set(null); - } - } catch (Exception e) { - log.error("[{}] Failed to fetch edge events", sessionId, e); - result.setException(e); - } - } - - private ConnectResponseMsg processConnect(ConnectRequestMsg request) { - log.trace("[{}] processConnect [{}]", sessionId, request); - Optional optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); - if (optional.isPresent()) { - edge = optional.get(); - tenantId = edge.getTenantId(); - try { - if (edge.getSecret().equals(request.getEdgeSecret())) { - sessionOpenListener.accept(edge.getId(), this); - edgeVersion = request.getEdgeVersion(); - processSaveEdgeVersionAsAttribute(request.getEdgeVersion().name()); - return ConnectResponseMsg.newBuilder() - .setResponseCode(ConnectResponseCode.ACCEPTED) - .setErrorMsg("") - .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)) - .setMaxInboundMessageSize(maxInboundMessageSize) - .build(); - } - String error = "Failed to validate the edge!"; - String failureMsg = String.format("%s Provided request secret: %s", error, request.getEdgeSecret()); - ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) - .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(error).build()); - return ConnectResponseMsg.newBuilder() - .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) - .setErrorMsg(failureMsg) - .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); - } catch (Exception e) { - String failureMsg = "Failed to process edge connection!"; - ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) - .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(e.getMessage()).build()); - log.error(failureMsg, e); - return ConnectResponseMsg.newBuilder() - .setResponseCode(ConnectResponseCode.SERVER_UNAVAILABLE) - .setErrorMsg(failureMsg) - .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); - } - } - return ConnectResponseMsg.newBuilder() - .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) - .setErrorMsg("Failed to find the edge! Routing key: " + request.getEdgeRoutingKey()) - .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); - } - - private void processSaveEdgeVersionAsAttribute(String edgeVersion) { - AttributeKvEntry attributeKvEntry = new BaseAttributeKvEntry(new StringDataEntry(DataConstants.EDGE_VERSION_ATTR_KEY, edgeVersion), System.currentTimeMillis()); - ctx.getAttributesService().save(tenantId, edge.getId(), AttributeScope.SERVER_SCOPE, attributeKvEntry); - } - - private void interruptGeneralProcessingOnSync() { - log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", tenantId, edge.getId(), sessionId); - stopCurrentSendDownlinkMsgsTask(true); - } - - protected ListenableFuture sendDownlinkMsgsPack(List downlinkMsgsPack) { - interruptPreviousSendDownlinkMsgsTask(); - - sessionState.setSendDownlinkMsgsFuture(SettableFuture.create()); - sessionState.getPendingMsgsMap().clear(); - - downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg)); - scheduleDownlinkMsgsPackSend(1); - - return sessionState.getSendDownlinkMsgsFuture(); - } - - private void interruptPreviousSendDownlinkMsgsTask() { - if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone() - || sessionState.getScheduledSendDownlinkTask() != null && !sessionState.getScheduledSendDownlinkTask().isCancelled()) { - log.debug("[{}][{}][{}] Previous send downlink future was not properly completed, stopping it now!", tenantId, edge.getId(), sessionId); - stopCurrentSendDownlinkMsgsTask(true); - } - } - - private void onUplinkMsg(UplinkMsg uplinkMsg) { - if (isRateLimitViolated(uplinkMsg)) { - return; - } - ListenableFuture> future = processUplinkMsg(uplinkMsg); - Futures.addCallback(future, new FutureCallback<>() { - @Override - public void onSuccess(@Nullable List result) { - sendResponseMessage(uplinkMsg.getUplinkMsgId(), true, null); - } - - @Override - public void onFailure(Throwable t) { - String errorMsg = EdgeUtils.createErrorMsgFromRootCauseAndStackTrace(t); - sendResponseMessage(uplinkMsg.getUplinkMsgId(), false, errorMsg); - } - }, ctx.getGrpcCallbackExecutorService()); - } - - private boolean isRateLimitViolated(UplinkMsg uplinkMsg) { - if (!ctx.getRateLimitService().checkRateLimit(LimitedApi.EDGE_UPLINK_MESSAGES, tenantId) || - !ctx.getRateLimitService().checkRateLimit(LimitedApi.EDGE_UPLINK_MESSAGES_PER_EDGE, tenantId, edge.getId())) { - String errorMsg = String.format("Failed to process uplink message. %s", RATE_LIMIT_REACHED); - sendResponseMessage(uplinkMsg.getUplinkMsgId(), false, errorMsg); - return true; - } - return false; - } - - private void scheduleDownlinkMsgsPackSend(int attempt) { - Runnable sendDownlinkMsgsTask = () -> { - try { - if (isConnected() && !sessionState.getPendingMsgsMap().values().isEmpty()) { - List copy = new ArrayList<>(sessionState.getPendingMsgsMap().values()); - if (attempt > 1) { - String error = "Failed to deliver the batch"; - String failureMsg = String.format("{%s}: {%s}", error, copy); - if (attempt == 2) { - // Send a failure notification only on the second attempt. - // This ensures that failure alerts are sent just once to avoid redundant notifications. - ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId) - .edgeId(edge.getId()).customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(error).build()); - } - log.warn("[{}][{}] {}, attempt: {}", tenantId, sessionId, failureMsg, attempt); - } - log.trace("[{}][{}][{}] downlink msg(s) are going to be send.", tenantId, sessionId, copy.size()); - for (DownlinkMsg downlinkMsg : copy) { - if (clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > clientMaxInboundMessageSize) { - String error = String.format("Client max inbound message size %s is exceeded. Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE " + - "env variable on the edge and restart it.", clientMaxInboundMessageSize); - String message = String.format("Downlink msg size %s exceeds client max inbound message size %s. " + - "Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE env variable on the edge and restart it.", downlinkMsg.getSerializedSize(), clientMaxInboundMessageSize); - log.error("[{}][{}][{}] {} Message {}", tenantId, edge.getId(), sessionId, message, downlinkMsg); - ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId) - .edgeId(edge.getId()).customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(message).error(error).build()); - sessionState.getPendingMsgsMap().remove(downlinkMsg.getDownlinkMsgId()); - } else { - sendDownlinkMsg(ResponseMsg.newBuilder() - .setDownlinkMsg(downlinkMsg) - .build()); - } - } - if (attempt < MAX_DOWNLINK_ATTEMPTS) { - scheduleDownlinkMsgsPackSend(attempt + 1); - } else { - String failureMsg = String.format("Failed to deliver messages: %s", copy); - log.warn("[{}][{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", - tenantId, sessionId, MAX_DOWNLINK_ATTEMPTS, copy); - ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) - .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg) - .error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build()); - stopCurrentSendDownlinkMsgsTask(false); - } - } else { - stopCurrentSendDownlinkMsgsTask(false); - } - } catch (Exception e) { - log.warn("[{}][{}] Failed to send downlink msgs. Error msg {}", tenantId, sessionId, e.getMessage(), e); - stopCurrentSendDownlinkMsgsTask(true); - } - }; - - if (attempt == 1) { - sendDownlinkExecutorService.submit(sendDownlinkMsgsTask); - } else { - sessionState.setScheduledSendDownlinkTask( - sendDownlinkExecutorService.schedule( - sendDownlinkMsgsTask, - ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches(), - TimeUnit.MILLISECONDS) - ); - } - } - - private void sendResponseMessage(int uplinkMsgId, boolean success, String errorMsg) { - UplinkResponseMsg.Builder responseBuilder = UplinkResponseMsg.newBuilder() - .setUplinkMsgId(uplinkMsgId) - .setSuccess(success); - if (errorMsg != null) { - responseBuilder.setErrorMsg(errorMsg); - } - sendDownlinkMsg(ResponseMsg.newBuilder() - .setUplinkResponseMsg(responseBuilder.build()) - .build()); - } - - private void onDownlinkResponse(DownlinkResponseMsg msg) { - try { - if (msg.getSuccess()) { - sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId()); - log.debug("[{}][{}] Msg has been processed successfully! Msg Id: [{}], Msg: {}", tenantId, edge.getRoutingKey(), msg.getDownlinkMsgId(), msg); - } else { - log.error("[{}][{}] Msg processing failed! Msg Id: [{}], Error msg: {}", tenantId, edge.getRoutingKey(), msg.getDownlinkMsgId(), msg.getErrorMsg()); - } - if (sessionState.getPendingMsgsMap().isEmpty()) { - log.debug("[{}][{}] Pending msgs map is empty. Stopping current iteration", tenantId, edge.getRoutingKey()); - stopCurrentSendDownlinkMsgsTask(false); - } - } catch (Exception e) { - log.error("[{}][{}] Can't process downlink response message [{}]", tenantId, sessionId, msg, e); - } - } - - @Override - public void processHighPriorityEvents() { - try { - List highPriorityEvents = new ArrayList<>(); - EdgeEvent event; - while ((event = highPriorityQueue.poll()) != null) { - highPriorityEvents.add(event); - } - List downlinkMsgsPack = convertToDownlinkMsgsPack(highPriorityEvents); - sendDownlinkMsgsPack(downlinkMsgsPack).get(); - } catch (Exception e) { - log.error("[{}] Failed to process high priority events", sessionId, e); - } - } - - @Override - public ListenableFuture processEdgeEvents() throws Exception { - SettableFuture result = SettableFuture.create(); - log.trace("[{}][{}] starting processing edge events", tenantId, sessionId); - if (isConnected() && isSyncCompleted()) { - Pair startTsAndSeqId = getQueueStartTsAndSeqId().get(); - previousStartTs = startTsAndSeqId.getFirst(); - previousStartSeqId = startTsAndSeqId.getSecond(); - GeneralEdgeEventFetcher fetcher = new GeneralEdgeEventFetcher( - previousStartTs, - previousStartSeqId, - seqIdEnd, - false, - Integer.toUnsignedLong(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()), - ctx.getEdgeEventService()); - Futures.addCallback(startProcessingEdgeEvents(fetcher), new FutureCallback<>() { - @Override - public void onSuccess(@Nullable Pair newStartTsAndSeqId) { - if (newStartTsAndSeqId != null) { - ListenableFuture> updateFuture = updateQueueStartTsAndSeqId(newStartTsAndSeqId); - Futures.addCallback(updateFuture, new FutureCallback<>() { - @Override - public void onSuccess(@Nullable List list) { - log.debug("[{}][{}] queue offset was updated [{}]", tenantId, sessionId, newStartTsAndSeqId); - if (fetcher.isSeqIdNewCycleStarted()) { - seqIdEnd = fetcher.getSeqIdEnd(); - boolean newEventsAvailable = isNewEdgeEventsAvailable(); - result.set(newEventsAvailable); - } else { - seqIdEnd = null; - boolean newEventsAvailable = isSeqIdStartedNewCycle(); - if (!newEventsAvailable) { - newEventsAvailable = isNewEdgeEventsAvailable(); - } - result.set(newEventsAvailable); - } - } - - @Override - public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to update queue offset [{}]", tenantId, sessionId, newStartTsAndSeqId, t); - result.setException(t); - } - }, ctx.getGrpcCallbackExecutorService()); - } else { - log.trace("[{}][{}] newStartTsAndSeqId is null. Skipping iteration without db update", tenantId, sessionId); - result.set(Boolean.FALSE); - } - } - - @Override - public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to process events", tenantId, sessionId, t); - result.setException(t); - } - }, ctx.getGrpcCallbackExecutorService()); - } else { - log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId); - result.set(null); - } - return result; - } - - protected List convertToDownlinkMsgsPack(List edgeEvents) { - List result = new ArrayList<>(); - for (EdgeEvent edgeEvent : edgeEvents) { - log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, sessionId, edgeEvent); - DownlinkMsg downlinkMsg = null; - try { - switch (edgeEvent.getAction()) { - case UPDATED, ADDED, DELETED, ASSIGNED_TO_EDGE, UNASSIGNED_FROM_EDGE, ALARM_ACK, ALARM_CLEAR, - ALARM_DELETE, CREDENTIALS_UPDATED, RELATION_ADD_OR_UPDATE, RELATION_DELETED, RPC_CALL, - ASSIGNED_TO_CUSTOMER, UNASSIGNED_FROM_CUSTOMER, ADDED_COMMENT, UPDATED_COMMENT, DELETED_COMMENT -> { - downlinkMsg = convertEntityEventToDownlink(edgeEvent); - if (downlinkMsg != null && downlinkMsg.getWidgetTypeUpdateMsgCount() > 0) { - log.trace("[{}][{}] widgetTypeUpdateMsg message processed, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsgId()); - } else { - log.trace("[{}][{}] entity message processed [{}]", tenantId, sessionId, downlinkMsg); - } - } - case ATTRIBUTES_UPDATED, POST_ATTRIBUTES, ATTRIBUTES_DELETED, TIMESERIES_UPDATED -> - downlinkMsg = ctx.getTelemetryProcessor().convertTelemetryEventToDownlink(edge, edgeEvent); - default -> log.warn("[{}][{}] Unsupported action type [{}]", tenantId, sessionId, edgeEvent.getAction()); - } - } catch (Exception e) { - log.error("[{}][{}] Exception during converting edge event to downlink msg", tenantId, sessionId, e); - } - if (downlinkMsg != null) { - result.add(downlinkMsg); - } - } - return result; - } - - private ListenableFuture> getQueueStartTsAndSeqId() { - ListenableFuture> future = - ctx.getAttributesService().find(edge.getTenantId(), edge.getId(), AttributeScope.SERVER_SCOPE, Arrays.asList(QUEUE_START_TS_ATTR_KEY, QUEUE_START_SEQ_ID_ATTR_KEY)); - return Futures.transform(future, attributeKvEntries -> { - long startTs = 0L; - long startSeqId = 0L; - for (AttributeKvEntry attributeKvEntry : attributeKvEntries) { - if (QUEUE_START_TS_ATTR_KEY.equals(attributeKvEntry.getKey())) { - startTs = attributeKvEntry.getLongValue().isPresent() ? attributeKvEntry.getLongValue().get() : 0L; - } - if (QUEUE_START_SEQ_ID_ATTR_KEY.equals(attributeKvEntry.getKey())) { - startSeqId = attributeKvEntry.getLongValue().isPresent() ? attributeKvEntry.getLongValue().get() : 0L; - } - } - if (startSeqId == 0L) { - startSeqId = findStartSeqIdFromOldestEventIfAny(); - } - return Pair.of(startTs, startSeqId); - }, ctx.getGrpcCallbackExecutorService()); - } - - private boolean isSeqIdStartedNewCycle() { - try { - TimePageLink pageLink = new TimePageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount(), 0, null, null, newStartTs, System.currentTimeMillis()); - PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), 0L, previousStartSeqId == 0 ? null : previousStartSeqId - 1, pageLink); - return !edgeEvents.getData().isEmpty(); - } catch (Exception e) { - log.error("[{}][{}][{}] Failed to execute isSeqIdStartedNewCycle", tenantId, edge.getId(), sessionId, e); - } - return false; - } - - private boolean isNewEdgeEventsAvailable() { - try { - TimePageLink pageLink = new TimePageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount(), 0, null, null, newStartTs, System.currentTimeMillis()); - PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), newStartSeqId, null, pageLink); - return !edgeEvents.getData().isEmpty() || !highPriorityQueue.isEmpty(); - } catch (Exception e) { - log.error("[{}][{}][{}] Failed to execute isNewEdgeEventsAvailable", tenantId, edge.getId(), sessionId, e); - } - return false; - } - - private long findStartSeqIdFromOldestEventIfAny() { - long startSeqId = 0L; - try { - TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("createdTime"), null, null); - PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), null, null, pageLink); - if (!edgeEvents.getData().isEmpty()) { - startSeqId = edgeEvents.getData().get(0).getSeqId() - 1; - } - } catch (Exception e) { - log.error("[{}][{}][{}] Failed to execute findStartSeqIdFromOldestEventIfAny", tenantId, edge.getId(), sessionId, e); - } - return startSeqId; - } - - private ListenableFuture> updateQueueStartTsAndSeqId(Pair pair) { - newStartTs = pair.getFirst(); - newStartSeqId = pair.getSecond(); - log.trace("[{}] updateQueueStartTsAndSeqId [{}][{}][{}]", sessionId, edge.getId(), newStartTs, newStartSeqId); - List attributes = Arrays.asList( - new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_TS_ATTR_KEY, newStartTs), System.currentTimeMillis()), - new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_SEQ_ID_ATTR_KEY, newStartSeqId), System.currentTimeMillis())); - return ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), AttributeScope.SERVER_SCOPE, attributes); - } - - protected ListenableFuture> startProcessingEdgeEvents(EdgeEventFetcher fetcher) { - SettableFuture> result = SettableFuture.create(); - PageLink pageLink = fetcher.getPageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()); - processEdgeEvents(fetcher, pageLink, result); - return result; - } - - private void markSyncCompletedSendEdgeEventUpdate() { - syncCompleted = true; - ctx.getClusterService().onEdgeEventUpdate(new EdgeEventUpdateMsg(edge.getTenantId(), edge.getId())); - } - - private void stopCurrentSendDownlinkMsgsTask(Boolean isInterrupted) { - if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) { - sessionState.getSendDownlinkMsgsFuture().set(isInterrupted); - } - if (sessionState.getScheduledSendDownlinkTask() != null) { - sessionState.getScheduledSendDownlinkTask().cancel(true); - } - } - - private void sendDownlinkMsg(ResponseMsg downlinkMsg) { - if (downlinkMsg.getDownlinkMsg().getWidgetTypeUpdateMsgCount() > 0) { - log.trace("[{}][{}] Sending downlink widgetTypeUpdateMsg, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); - } else { - log.trace("[{}][{}] Sending downlink msg [{}]", tenantId, sessionId, downlinkMsg); - } - if (isConnected()) { - downlinkMsgLock.lock(); - try { - outputStream.onNext(downlinkMsg); - } catch (Exception e) { - log.error("[{}][{}] Failed to send downlink message [{}]", tenantId, sessionId, downlinkMsg, e); - connected = false; - sessionCloseListener.accept(edge, sessionId); - } finally { - downlinkMsgLock.unlock(); - } - log.trace("[{}][{}] Response msg successfully sent. downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); - } - } - - protected DownlinkMsg convertEntityEventToDownlink(EdgeEvent edgeEvent) { - log.trace("[{}] Executing convertEntityEventToDownlink, edgeEvent [{}], action [{}]", edgeEvent.getTenantId(), edgeEvent, edgeEvent.getAction()); - return switch (edgeEvent.getType()) { - case EDGE -> ctx.getEdgeProcessor().convertEdgeEventToDownlink(edgeEvent); - case DEVICE -> ctx.getDeviceProcessor().convertDeviceEventToDownlink(edgeEvent, edgeVersion); - case DEVICE_PROFILE -> ctx.getDeviceProfileProcessor().convertDeviceProfileEventToDownlink(edgeEvent, edgeVersion); - case ASSET_PROFILE -> ctx.getAssetProfileProcessor().convertAssetProfileEventToDownlink(edgeEvent, edgeVersion); - case ASSET -> ctx.getAssetProcessor().convertAssetEventToDownlink(edgeEvent, edgeVersion); - case ENTITY_VIEW -> ctx.getEntityViewProcessor().convertEntityViewEventToDownlink(edgeEvent, edgeVersion); - case DASHBOARD -> ctx.getDashboardProcessor().convertDashboardEventToDownlink(edgeEvent, edgeVersion); - case CUSTOMER -> ctx.getCustomerProcessor().convertCustomerEventToDownlink(edgeEvent, edgeVersion); - case RULE_CHAIN -> ctx.getRuleChainProcessor().convertRuleChainEventToDownlink(edgeEvent, edgeVersion); - case RULE_CHAIN_METADATA -> ctx.getRuleChainProcessor().convertRuleChainMetadataEventToDownlink(edgeEvent, edgeVersion); - case ALARM -> ctx.getAlarmProcessor().convertAlarmEventToDownlink(edgeEvent, edgeVersion); - case ALARM_COMMENT -> ctx.getAlarmProcessor().convertAlarmCommentEventToDownlink(edgeEvent, edgeVersion); - case USER -> ctx.getUserProcessor().convertUserEventToDownlink(edgeEvent, edgeVersion); - case RELATION -> ctx.getRelationProcessor().convertRelationEventToDownlink(edgeEvent, edgeVersion); - case WIDGETS_BUNDLE -> ctx.getWidgetBundleProcessor().convertWidgetsBundleEventToDownlink(edgeEvent, edgeVersion); - case WIDGET_TYPE -> ctx.getWidgetTypeProcessor().convertWidgetTypeEventToDownlink(edgeEvent, edgeVersion); - case ADMIN_SETTINGS -> ctx.getAdminSettingsProcessor().convertAdminSettingsEventToDownlink(edgeEvent, edgeVersion); - case OTA_PACKAGE -> ctx.getOtaPackageProcessor().convertOtaPackageEventToDownlink(edgeEvent, edgeVersion); - case TB_RESOURCE -> ctx.getResourceProcessor().convertResourceEventToDownlink(edgeEvent, edgeVersion); - case QUEUE -> ctx.getQueueProcessor().convertQueueEventToDownlink(edgeEvent, edgeVersion); - case TENANT -> ctx.getTenantProcessor().convertTenantEventToDownlink(edgeEvent, edgeVersion); - case TENANT_PROFILE -> ctx.getTenantProfileProcessor().convertTenantProfileEventToDownlink(edgeEvent, edgeVersion); - case NOTIFICATION_RULE -> ctx.getNotificationEdgeProcessor().convertNotificationRuleToDownlink(edgeEvent); - case NOTIFICATION_TARGET -> ctx.getNotificationEdgeProcessor().convertNotificationTargetToDownlink(edgeEvent); - case NOTIFICATION_TEMPLATE -> ctx.getNotificationEdgeProcessor().convertNotificationTemplateToDownlink(edgeEvent); - case OAUTH2_CLIENT -> ctx.getOAuth2EdgeProcessor().convertOAuth2ClientEventToDownlink(edgeEvent, edgeVersion); - case DOMAIN -> ctx.getOAuth2EdgeProcessor().convertOAuth2DomainEventToDownlink(edgeEvent, edgeVersion); - default -> { - log.warn("[{}] Unsupported edge event type [{}]", edgeEvent.getTenantId(), edgeEvent); - yield null; - } - }; - } - - public void addEventToHighPriorityQueue(EdgeEvent edgeEvent) { - while (highPriorityQueue.size() > maxHighPriorityQueueSizePerSession) { - EdgeEvent oldestHighPriority = highPriorityQueue.poll(); - if (oldestHighPriority != null) { - log.warn("[{}][{}][{}] High priority queue is full. Removing oldest high priority event from queue {}", - tenantId, edge.getId(), sessionId, oldestHighPriority); - } - } - highPriorityQueue.add(edgeEvent); - } - - protected ListenableFuture> processUplinkMsg(UplinkMsg uplinkMsg) { - List> result = new ArrayList<>(); - try { - if (uplinkMsg.getDeviceProfileUpdateMsgCount() > 0) { - for (DeviceProfileUpdateMsg deviceProfileUpdateMsg : uplinkMsg.getDeviceProfileUpdateMsgList()) { - result.add(((DeviceProfileProcessor) ctx.getDeviceProfileEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processDeviceProfileMsgFromEdge(edge.getTenantId(), edge, deviceProfileUpdateMsg)); - } - } - if (uplinkMsg.getDeviceUpdateMsgCount() > 0) { - for (DeviceUpdateMsg deviceUpdateMsg : uplinkMsg.getDeviceUpdateMsgList()) { - result.add(((DeviceProcessor) ctx.getDeviceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processDeviceMsgFromEdge(edge.getTenantId(), edge, deviceUpdateMsg)); - } - } - if (uplinkMsg.getDeviceCredentialsUpdateMsgCount() > 0) { - for (DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg : uplinkMsg.getDeviceCredentialsUpdateMsgList()) { - result.add(((DeviceProcessor) ctx.getDeviceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processDeviceCredentialsMsgFromEdge(edge.getTenantId(), edge.getId(), deviceCredentialsUpdateMsg)); - } - } - if (uplinkMsg.getAssetProfileUpdateMsgCount() > 0) { - for (AssetProfileUpdateMsg assetProfileUpdateMsg : uplinkMsg.getAssetProfileUpdateMsgList()) { - result.add(((AssetProfileProcessor) ctx.getAssetProfileEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processAssetProfileMsgFromEdge(edge.getTenantId(), edge, assetProfileUpdateMsg)); - } - } - if (uplinkMsg.getAssetUpdateMsgCount() > 0) { - for (AssetUpdateMsg assetUpdateMsg : uplinkMsg.getAssetUpdateMsgList()) { - result.add(((AssetProcessor) ctx.getAssetEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processAssetMsgFromEdge(edge.getTenantId(), edge, assetUpdateMsg)); - } - } - if (uplinkMsg.getEntityViewUpdateMsgCount() > 0) { - for (EntityViewUpdateMsg entityViewUpdateMsg : uplinkMsg.getEntityViewUpdateMsgList()) { - result.add(((EntityViewProcessor) ctx.getEntityViewProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processEntityViewMsgFromEdge(edge.getTenantId(), edge, entityViewUpdateMsg)); - } - } - if (uplinkMsg.getEntityDataCount() > 0) { - for (EntityDataProto entityData : uplinkMsg.getEntityDataList()) { - result.addAll(ctx.getTelemetryProcessor().processTelemetryMsg(edge.getTenantId(), entityData)); - } - } - if (uplinkMsg.getAlarmUpdateMsgCount() > 0) { - for (AlarmUpdateMsg alarmUpdateMsg : uplinkMsg.getAlarmUpdateMsgList()) { - result.add(((AlarmProcessor) ctx.getAlarmEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processAlarmMsgFromEdge(edge.getTenantId(), edge.getId(), alarmUpdateMsg)); - } - } - if (uplinkMsg.getAlarmCommentUpdateMsgCount() > 0) { - for (AlarmCommentUpdateMsg alarmCommentUpdateMsg : uplinkMsg.getAlarmCommentUpdateMsgList()) { - result.add(((AlarmProcessor) ctx.getAlarmEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processAlarmCommentMsgFromEdge(edge.getTenantId(), edge.getId(), alarmCommentUpdateMsg)); - } - } - if (uplinkMsg.getRelationUpdateMsgCount() > 0) { - for (RelationUpdateMsg relationUpdateMsg : uplinkMsg.getRelationUpdateMsgList()) { - result.add(((RelationProcessor) ctx.getRelationEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processRelationMsgFromEdge(edge.getTenantId(), edge, relationUpdateMsg)); - } - } - if (uplinkMsg.getDashboardUpdateMsgCount() > 0) { - for (DashboardUpdateMsg dashboardUpdateMsg : uplinkMsg.getDashboardUpdateMsgList()) { - result.add(((DashboardProcessor) ctx.getDashboardEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processDashboardMsgFromEdge(edge.getTenantId(), edge, dashboardUpdateMsg)); - } - } - if (uplinkMsg.getResourceUpdateMsgCount() > 0) { - for (ResourceUpdateMsg resourceUpdateMsg : uplinkMsg.getResourceUpdateMsgList()) { - result.add(((ResourceProcessor) ctx.getResourceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processResourceMsgFromEdge(edge.getTenantId(), edge, resourceUpdateMsg)); - } - } - if (uplinkMsg.getRuleChainMetadataRequestMsgCount() > 0) { - for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processRuleChainMetadataRequestMsg(edge.getTenantId(), edge, ruleChainMetadataRequestMsg)); - } - } - if (uplinkMsg.getAttributesRequestMsgCount() > 0) { - for (AttributesRequestMsg attributesRequestMsg : uplinkMsg.getAttributesRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processAttributesRequestMsg(edge.getTenantId(), edge, attributesRequestMsg)); - } - } - if (uplinkMsg.getRelationRequestMsgCount() > 0) { - for (RelationRequestMsg relationRequestMsg : uplinkMsg.getRelationRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processRelationRequestMsg(edge.getTenantId(), edge, relationRequestMsg)); - } - } - if (uplinkMsg.getUserCredentialsRequestMsgCount() > 0) { - for (UserCredentialsRequestMsg userCredentialsRequestMsg : uplinkMsg.getUserCredentialsRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processUserCredentialsRequestMsg(edge.getTenantId(), edge, userCredentialsRequestMsg)); - } - } - if (uplinkMsg.getDeviceCredentialsRequestMsgCount() > 0) { - for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg : uplinkMsg.getDeviceCredentialsRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processDeviceCredentialsRequestMsg(edge.getTenantId(), edge, deviceCredentialsRequestMsg)); - } - } - if (uplinkMsg.getDeviceRpcCallMsgCount() > 0) { - for (DeviceRpcCallMsg deviceRpcCallMsg : uplinkMsg.getDeviceRpcCallMsgList()) { - result.add(((DeviceProcessor) ctx.getDeviceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) - .processDeviceRpcCallFromEdge(edge.getTenantId(), edge, deviceRpcCallMsg)); - } - } - if (uplinkMsg.getWidgetBundleTypesRequestMsgCount() > 0) { - for (WidgetBundleTypesRequestMsg widgetBundleTypesRequestMsg : uplinkMsg.getWidgetBundleTypesRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processWidgetBundleTypesRequestMsg(edge.getTenantId(), edge, widgetBundleTypesRequestMsg)); - } - } - if (uplinkMsg.getEntityViewsRequestMsgCount() > 0) { - for (EntityViewsRequestMsg entityViewRequestMsg : uplinkMsg.getEntityViewsRequestMsgList()) { - result.add(ctx.getEdgeRequestsService().processEntityViewsRequestMsg(edge.getTenantId(), edge, entityViewRequestMsg)); - } - } - } catch (Exception e) { - String failureMsg = String.format("Can't process uplink msg [%s] from edge", uplinkMsg); - log.error("[{}][{}] Can't process uplink msg [{}]", edge.getTenantId(), sessionId, uplinkMsg, e); - ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(edge.getTenantId()).edgeId(edge.getId()) - .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(e.getMessage()).build()); - return Futures.immediateFailedFuture(e); - } - return Futures.allAsList(result); - } - - @Override - public void close() { - log.debug("[{}][{}] Closing session", tenantId, sessionId); - connected = false; - try { - outputStream.onCompleted(); - } catch (Exception e) { - log.debug("[{}][{}] Failed to close output stream: {}", tenantId, sessionId, e.getMessage()); - } - } - -} 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 1b1a6de088..9c0aa101c4 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 @@ -58,6 +58,9 @@ import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc; import org.thingsboard.server.gen.edge.v1.RequestMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; +import org.thingsboard.server.queue.kafka.TbKafkaSettings; +import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeContextComponent; @@ -69,6 +72,7 @@ import java.io.InputStream; import java.util.Collections; import java.util.HashMap; import java.util.Map; +import java.util.Optional; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -86,7 +90,7 @@ import java.util.function.Consumer; @TbCoreComponent public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { - private final ConcurrentMap sessions = new ConcurrentHashMap<>(); + private final ConcurrentMap sessions = new ConcurrentHashMap<>(); private final ConcurrentMap sessionNewEventsLocks = new ConcurrentHashMap<>(); private final Map sessionNewEvents = new HashMap<>(); private final ConcurrentMap> sessionEdgeEventChecks = new ConcurrentHashMap<>(); @@ -120,9 +124,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Value("${edges.max_high_priority_queue_size_per_session:10000}") private int maxHighPriorityQueueSizePerSession; - @Value("#{'${queue.type:null}' == 'kafka'}") - private boolean isKafkaSupported; - @Autowired @Lazy private EdgeContextComponent ctx; @@ -139,9 +140,18 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Autowired private TbTransactionalCache edgeIdServiceIdCache; + @Autowired + private TopicService topicService; + @Autowired private TbCoreQueueFactory tbCoreQueueFactory; + @Autowired + private Optional kafkaSettings; + + @Autowired + private Optional kafkaTopicConfigs; + private Server server; private ScheduledExecutorService edgeEventProcessingExecutorService; @@ -210,13 +220,13 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public StreamObserver handleMsgs(StreamObserver outputStream) { - AbstractEdgeGrpcSession session = createEdgeGrpcSession(outputStream); + EdgeGrpcSession session = createEdgeGrpcSession(outputStream); return session.getInputStream(); } - private AbstractEdgeGrpcSession createEdgeGrpcSession(StreamObserver outputStream) { - return isKafkaSupported - ? new KafkaEdgeGrpcSession(ctx, tbCoreQueueFactory, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, + private EdgeGrpcSession createEdgeGrpcSession(StreamObserver outputStream) { + return kafkaSettings.isPresent() && kafkaTopicConfigs.isPresent() + ? new KafkaEdgeGrpcSession(ctx, topicService, tbCoreQueueFactory, kafkaSettings.get(), kafkaTopicConfigs.get(), outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession) : new PostgresEdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); @@ -250,7 +260,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void updateEdge(TenantId tenantId, Edge edge) { - AbstractEdgeGrpcSession session = sessions.get(edge.getId()); + EdgeGrpcSession session = sessions.get(edge.getId()); if (session != null && session.isConnected()) { log.debug("[{}] Updating configuration for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); session.onConfigurationUpdate(edge); @@ -261,9 +271,11 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void deleteEdge(TenantId tenantId, EdgeId edgeId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + EdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); + session.destroy(); + session.deleteTopic(edgeId); session.close(); sessions.remove(edgeId); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); @@ -278,7 +290,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + EdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEventUpdate [{}]", tenantId, edgeId.getId()); updateSessionEventsFlag(tenantId, edgeId); @@ -289,7 +301,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i TenantId tenantId = msg.getTenantId(); EdgeEvent edgeEvent = msg.getEdgeEvent(); EdgeId edgeId = edgeEvent.getEdgeId(); - AbstractEdgeGrpcSession session = sessions.get(edgeId); + EdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId); session.addEventToHighPriorityQueue(edgeEvent); @@ -310,7 +322,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void onEdgeConnect(EdgeId edgeId, AbstractEdgeGrpcSession edgeGrpcSession) { + private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { Edge edge = edgeGrpcSession.getEdge(); TenantId tenantId = edge.getTenantId(); log.info("[{}][{}] edge [{}] connected successfully.", tenantId, edgeGrpcSession.getSessionId(), edgeId); @@ -333,7 +345,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId, String requestServiceId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + EdgeGrpcSession session = sessions.get(edgeId); if (session != null) { if (!session.isSyncCompleted()) { clusterService.pushEdgeSyncResponseToCore(new FromEdgeSyncResponse(requestId, tenantId, edgeId, false, "Sync process is active at the moment"), requestServiceId); @@ -353,7 +365,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i ToEdgeSyncRequest request = new ToEdgeSyncRequest(UUID.randomUUID(), tenantId, edgeId, serviceInfoProvider.getServiceId()); UUID requestId = request.getId(); - AbstractEdgeGrpcSession session = sessions.get(request.getEdgeId()); + EdgeGrpcSession session = sessions.get(request.getEdgeId()); if (session != null && !session.isSyncCompleted()) { responseConsumer.accept(new FromEdgeSyncResponse(requestId, request.getTenantId(), request.getEdgeId(), false, "Sync process is active at the moment")); } else { @@ -387,7 +399,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void scheduleEdgeEventsCheck(AbstractEdgeGrpcSession session) { + private void scheduleEdgeEventsCheck(EdgeGrpcSession session) { EdgeId edgeId = session.getEdge().getId(); TenantId tenantId = session.getEdge().getTenantId(); if (sessions.containsKey(edgeId)) { @@ -434,7 +446,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void processEdgeEventMigrationIfNeeded(AbstractEdgeGrpcSession session, EdgeId edgeId) throws Exception { + private void processEdgeEventMigrationIfNeeded(EdgeGrpcSession session, EdgeId edgeId) throws Exception { boolean isMigrationProcessed = edgeEventsProcessed.getOrDefault(edgeId, Boolean.FALSE); if (!isMigrationProcessed) { Boolean eventsExist = session.migrateEdgeEvents().get(); @@ -463,7 +475,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void onEdgeDisconnect(Edge edge, UUID sessionId) { EdgeId edgeId = edge.getId(); log.info("[{}][{}] edge disconnected!", edgeId, sessionId); - AbstractEdgeGrpcSession toRemove = sessions.get(edgeId); + EdgeGrpcSession toRemove = sessions.get(edgeId); if (toRemove.getSessionId().equals(sessionId)) { toRemove = sessions.remove(edgeId); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); 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 880712d6b2..4ad9c04353 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 @@ -15,23 +15,931 @@ */ package org.thingsboard.server.service.edge.rpc; +import com.datastax.oss.driver.api.core.uuid.Uuids; +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.SettableFuture; +import io.grpc.stub.StreamObserver; +import lombok.Data; +import lombok.extern.slf4j.Slf4j; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.springframework.data.util.Pair; +import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.DataConstants; +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.id.EdgeId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.limit.LimitedApi; +import org.thingsboard.server.common.data.notification.rule.trigger.EdgeCommunicationFailureTrigger; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.page.SortOrder; +import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg; +import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; +import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; +import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; +import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; +import org.thingsboard.server.gen.edge.v1.AttributesRequestMsg; +import org.thingsboard.server.gen.edge.v1.ConnectRequestMsg; +import org.thingsboard.server.gen.edge.v1.ConnectResponseCode; +import org.thingsboard.server.gen.edge.v1.ConnectResponseMsg; +import org.thingsboard.server.gen.edge.v1.DashboardUpdateMsg; +import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; +import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.v1.DeviceProfileUpdateMsg; +import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg; +import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; +import org.thingsboard.server.gen.edge.v1.DownlinkMsg; +import org.thingsboard.server.gen.edge.v1.DownlinkResponseMsg; +import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; +import org.thingsboard.server.gen.edge.v1.EdgeUpdateMsg; +import org.thingsboard.server.gen.edge.v1.EdgeVersion; +import org.thingsboard.server.gen.edge.v1.EntityDataProto; +import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; +import org.thingsboard.server.gen.edge.v1.EntityViewsRequestMsg; +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; +import org.thingsboard.server.gen.edge.v1.UplinkMsg; +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.EdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.GeneralEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.processor.alarm.AlarmProcessor; +import org.thingsboard.server.service.edge.rpc.processor.asset.AssetProcessor; +import org.thingsboard.server.service.edge.rpc.processor.asset.profile.AssetProfileProcessor; +import org.thingsboard.server.service.edge.rpc.processor.dashboard.DashboardProcessor; +import org.thingsboard.server.service.edge.rpc.processor.device.DeviceProcessor; +import org.thingsboard.server.service.edge.rpc.processor.device.profile.DeviceProfileProcessor; +import org.thingsboard.server.service.edge.rpc.processor.entityview.EntityViewProcessor; +import org.thingsboard.server.service.edge.rpc.processor.relation.RelationProcessor; +import org.thingsboard.server.service.edge.rpc.processor.resource.ResourceProcessor; -public interface EdgeGrpcSession { +import java.io.Closeable; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.BiConsumer; - void onConfigurationUpdate(Edge edge); +@Slf4j +@Data +public abstract class EdgeGrpcSession implements Closeable { - void startSyncProcess(boolean fullSync); + private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; + private static final String QUEUE_START_SEQ_ID_ATTR_KEY = "queueStartSeqId"; - boolean isConnected(); + private static final int MAX_DOWNLINK_ATTEMPTS = 10; + private static final String RATE_LIMIT_REACHED = "Rate limit reached"; - void destroy(); + protected static final ConcurrentLinkedQueue highPriorityQueue = new ConcurrentLinkedQueue<>(); - ListenableFuture migrateEdgeEvents() throws Exception; + protected UUID sessionId; + private BiConsumer sessionOpenListener; + private BiConsumer sessionCloseListener; - ListenableFuture processEdgeEvents() throws Exception; + private final EdgeSessionState sessionState = new EdgeSessionState(); + private final ReentrantLock downlinkMsgLock = new ReentrantLock(); - void processHighPriorityEvents(); + protected EdgeContextComponent ctx; + protected Edge edge; + protected TenantId tenantId; + + private Long newStartTs; + private Long previousStartTs; + private Long newStartSeqId; + private Long previousStartSeqId; + private Long seqIdEnd; + + private StreamObserver inputStream; + private StreamObserver outputStream; + + private boolean connected; + private volatile boolean syncCompleted; + + private EdgeVersion edgeVersion; + private int maxInboundMessageSize; + private int clientMaxInboundMessageSize; + private int maxHighPriorityQueueSizePerSession; + + private ScheduledExecutorService sendDownlinkExecutorService; + + public EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, + BiConsumer sessionOpenListener, + BiConsumer sessionCloseListener, + ScheduledExecutorService sendDownlinkExecutorService, + int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { + this.sessionId = UUID.randomUUID(); + this.ctx = ctx; + this.outputStream = outputStream; + this.sessionOpenListener = sessionOpenListener; + this.sessionCloseListener = sessionCloseListener; + this.sendDownlinkExecutorService = sendDownlinkExecutorService; + this.maxInboundMessageSize = maxInboundMessageSize; + this.maxHighPriorityQueueSizePerSession = maxHighPriorityQueueSizePerSession; + initInputStream(); + } + + protected abstract ListenableFuture migrateEdgeEvents() throws Exception; + + public void initInputStream() { + inputStream = new StreamObserver<>() { + @Override + public void onNext(RequestMsg requestMsg) { + if (!connected && requestMsg.getMsgType().equals(RequestMsgType.CONNECT_RPC_MESSAGE)) { + ConnectResponseMsg responseMsg = processConnect(requestMsg.getConnectRequestMsg()); + outputStream.onNext(ResponseMsg.newBuilder() + .setConnectResponseMsg(responseMsg) + .build()); + if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) { + outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); + } else { + if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) { + log.debug("[{}][{}] Client max inbound message size: {}", tenantId, sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); + clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); + } + connected = true; + } + } + if (connected) { + if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { + if (requestMsg.hasSyncRequestMsg()) { + boolean fullSync = false; + if (requestMsg.getSyncRequestMsg().hasFullSync()) { + fullSync = requestMsg.getSyncRequestMsg().getFullSync(); + } + startSyncProcess(fullSync); + } else { + syncCompleted = true; + } + } + if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) { + if (requestMsg.hasUplinkMsg()) { + onUplinkMsg(requestMsg.getUplinkMsg()); + } + if (requestMsg.hasDownlinkResponseMsg()) { + onDownlinkResponse(requestMsg.getDownlinkResponseMsg()); + } + } + } + } + + @Override + public void onError(Throwable t) { + log.error("[{}][{}] Stream was terminated due to error:", tenantId, sessionId, t); + closeSession(); + } + + @Override + public void onCompleted() { + log.info("[{}][{}] Stream was closed and completed successfully!", tenantId, sessionId); + closeSession(); + } + + private void closeSession() { + connected = false; + if (edge != null) { + try { + sessionCloseListener.accept(edge, sessionId); + } catch (Exception ignored) { + } + } + try { + outputStream.onCompleted(); + } catch (Exception ignored) { + } + } + }; + } + + public void onConfigurationUpdate(Edge edge) { + log.debug("[{}] onConfigurationUpdate [{}]", sessionId, edge); + this.tenantId = edge.getTenantId(); + this.edge = edge; + EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() + .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)).build(); + ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder() + .setEdgeUpdateMsg(edgeConfig) + .build(); + sendDownlinkMsg(edgeConfigMsg); + } + + public void startSyncProcess(boolean fullSync) { + log.info("[{}][{}][{}] Staring edge sync process", tenantId, edge.getId(), sessionId); + syncCompleted = false; + interruptGeneralProcessingOnSync(); + doSync(new EdgeSyncCursor(ctx, edge, fullSync)); + } + + private void doSync(EdgeSyncCursor cursor) { + if (cursor.hasNext()) { + EdgeEventFetcher next = cursor.getNext(); + log.info("[{}][{}] starting sync process, cursor current idx = {}, class = {}", + tenantId, edge.getId(), cursor.getCurrentIdx(), next.getClass().getSimpleName()); + ListenableFuture> future = startProcessingEdgeEvents(next); + Futures.addCallback(future, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Pair result) { + doSync(cursor); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Exception during sync process", tenantId, edge.getId(), t); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + log.info("[{}][{}] sync process completed", tenantId, edge.getId()); + DownlinkMsg syncCompleteDownlinkMsg = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .setSyncCompletedMsg(SyncCompletedMsg.newBuilder().build()) + .build(); + Futures.addCallback(sendDownlinkMsgsPack(Collections.singletonList(syncCompleteDownlinkMsg)), new FutureCallback<>() { + @Override + public void onSuccess(Boolean isInterrupted) { + markSyncCompletedSendEdgeEventUpdate(); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Exception during sending sync complete", tenantId, edge.getId(), t); + markSyncCompletedSendEdgeEventUpdate(); + } + }, ctx.getGrpcCallbackExecutorService()); + } + } + + protected void processEdgeEvents(EdgeEventFetcher fetcher, PageLink pageLink, SettableFuture> result) { + try { + if (!highPriorityQueue.isEmpty()) { + processHighPriorityEvents(); + } + PageData pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); + if (isConnected() && !pageData.getData().isEmpty()) { + log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, sessionId, pageData.getData().size()); + List downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); + Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Boolean isInterrupted) { + if (Boolean.TRUE.equals(isInterrupted)) { + log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId); + result.set(null); + } else { + if (isConnected() && pageData.hasNext()) { + processEdgeEvents(fetcher, pageLink.nextPageLink(), result); + } else { + EdgeEvent latestEdgeEvent = pageData.getData().get(pageData.getData().size() - 1); + UUID idOffset = latestEdgeEvent.getUuidId(); + if (idOffset != null) { + Long newStartTs = Uuids.unixTimestamp(idOffset); + long newStartSeqId = latestEdgeEvent.getSeqId(); + result.set(Pair.of(newStartTs, newStartSeqId)); + } else { + result.set(null); + } + } + } + } + + @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", sessionId); + result.set(null); + } + } catch (Exception e) { + log.error("[{}] Failed to fetch edge events", sessionId, e); + result.setException(e); + } + } + + private ConnectResponseMsg processConnect(ConnectRequestMsg request) { + log.trace("[{}] processConnect [{}]", sessionId, request); + Optional optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); + if (optional.isPresent()) { + edge = optional.get(); + tenantId = edge.getTenantId(); + try { + if (edge.getSecret().equals(request.getEdgeSecret())) { + sessionOpenListener.accept(edge.getId(), this); + edgeVersion = request.getEdgeVersion(); + processSaveEdgeVersionAsAttribute(request.getEdgeVersion().name()); + return ConnectResponseMsg.newBuilder() + .setResponseCode(ConnectResponseCode.ACCEPTED) + .setErrorMsg("") + .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)) + .setMaxInboundMessageSize(maxInboundMessageSize) + .build(); + } + String error = "Failed to validate the edge!"; + String failureMsg = String.format("%s Provided request secret: %s", error, request.getEdgeSecret()); + ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) + .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(error).build()); + return ConnectResponseMsg.newBuilder() + .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) + .setErrorMsg(failureMsg) + .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); + } catch (Exception e) { + String failureMsg = "Failed to process edge connection!"; + ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) + .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(e.getMessage()).build()); + log.error(failureMsg, e); + return ConnectResponseMsg.newBuilder() + .setResponseCode(ConnectResponseCode.SERVER_UNAVAILABLE) + .setErrorMsg(failureMsg) + .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); + } + } + return ConnectResponseMsg.newBuilder() + .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) + .setErrorMsg("Failed to find the edge! Routing key: " + request.getEdgeRoutingKey()) + .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); + } + + private void processSaveEdgeVersionAsAttribute(String edgeVersion) { + AttributeKvEntry attributeKvEntry = new BaseAttributeKvEntry(new StringDataEntry(DataConstants.EDGE_VERSION_ATTR_KEY, edgeVersion), System.currentTimeMillis()); + ctx.getAttributesService().save(tenantId, edge.getId(), AttributeScope.SERVER_SCOPE, attributeKvEntry); + } + + private void interruptGeneralProcessingOnSync() { + log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", tenantId, edge.getId(), sessionId); + stopCurrentSendDownlinkMsgsTask(true); + } + + protected ListenableFuture sendDownlinkMsgsPack(List downlinkMsgsPack) { + interruptPreviousSendDownlinkMsgsTask(); + + sessionState.setSendDownlinkMsgsFuture(SettableFuture.create()); + sessionState.getPendingMsgsMap().clear(); + + downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg)); + scheduleDownlinkMsgsPackSend(1); + + return sessionState.getSendDownlinkMsgsFuture(); + } + + private void interruptPreviousSendDownlinkMsgsTask() { + if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone() + || sessionState.getScheduledSendDownlinkTask() != null && !sessionState.getScheduledSendDownlinkTask().isCancelled()) { + log.debug("[{}][{}][{}] Previous send downlink future was not properly completed, stopping it now!", tenantId, edge.getId(), sessionId); + stopCurrentSendDownlinkMsgsTask(true); + } + } + + private void onUplinkMsg(UplinkMsg uplinkMsg) { + if (isRateLimitViolated(uplinkMsg)) { + return; + } + ListenableFuture> future = processUplinkMsg(uplinkMsg); + Futures.addCallback(future, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable List result) { + sendResponseMessage(uplinkMsg.getUplinkMsgId(), true, null); + } + + @Override + public void onFailure(Throwable t) { + String errorMsg = EdgeUtils.createErrorMsgFromRootCauseAndStackTrace(t); + sendResponseMessage(uplinkMsg.getUplinkMsgId(), false, errorMsg); + } + }, ctx.getGrpcCallbackExecutorService()); + } + + private boolean isRateLimitViolated(UplinkMsg uplinkMsg) { + if (!ctx.getRateLimitService().checkRateLimit(LimitedApi.EDGE_UPLINK_MESSAGES, tenantId) || + !ctx.getRateLimitService().checkRateLimit(LimitedApi.EDGE_UPLINK_MESSAGES_PER_EDGE, tenantId, edge.getId())) { + String errorMsg = String.format("Failed to process uplink message. %s", RATE_LIMIT_REACHED); + sendResponseMessage(uplinkMsg.getUplinkMsgId(), false, errorMsg); + return true; + } + return false; + } + + private void scheduleDownlinkMsgsPackSend(int attempt) { + Runnable sendDownlinkMsgsTask = () -> { + try { + if (isConnected() && !sessionState.getPendingMsgsMap().values().isEmpty()) { + List copy = new ArrayList<>(sessionState.getPendingMsgsMap().values()); + if (attempt > 1) { + String error = "Failed to deliver the batch"; + String failureMsg = String.format("{%s}: {%s}", error, copy); + if (attempt == 2) { + // Send a failure notification only on the second attempt. + // This ensures that failure alerts are sent just once to avoid redundant notifications. + ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId) + .edgeId(edge.getId()).customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(error).build()); + } + log.warn("[{}][{}] {}, attempt: {}", tenantId, sessionId, failureMsg, attempt); + } + log.trace("[{}][{}][{}] downlink msg(s) are going to be send.", tenantId, sessionId, copy.size()); + for (DownlinkMsg downlinkMsg : copy) { + if (clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > clientMaxInboundMessageSize) { + String error = String.format("Client max inbound message size %s is exceeded. Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE " + + "env variable on the edge and restart it.", clientMaxInboundMessageSize); + String message = String.format("Downlink msg size %s exceeds client max inbound message size %s. " + + "Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE env variable on the edge and restart it.", downlinkMsg.getSerializedSize(), clientMaxInboundMessageSize); + log.error("[{}][{}][{}] {} Message {}", tenantId, edge.getId(), sessionId, message, downlinkMsg); + ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId) + .edgeId(edge.getId()).customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(message).error(error).build()); + sessionState.getPendingMsgsMap().remove(downlinkMsg.getDownlinkMsgId()); + } else { + sendDownlinkMsg(ResponseMsg.newBuilder() + .setDownlinkMsg(downlinkMsg) + .build()); + } + } + if (attempt < MAX_DOWNLINK_ATTEMPTS) { + scheduleDownlinkMsgsPackSend(attempt + 1); + } else { + String failureMsg = String.format("Failed to deliver messages: %s", copy); + log.warn("[{}][{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", + tenantId, sessionId, MAX_DOWNLINK_ATTEMPTS, copy); + ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) + .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg) + .error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build()); + stopCurrentSendDownlinkMsgsTask(false); + } + } else { + stopCurrentSendDownlinkMsgsTask(false); + } + } catch (Exception e) { + log.warn("[{}][{}] Failed to send downlink msgs. Error msg {}", tenantId, sessionId, e.getMessage(), e); + stopCurrentSendDownlinkMsgsTask(true); + } + }; + + if (attempt == 1) { + sendDownlinkExecutorService.submit(sendDownlinkMsgsTask); + } else { + sessionState.setScheduledSendDownlinkTask( + sendDownlinkExecutorService.schedule( + sendDownlinkMsgsTask, + ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches(), + TimeUnit.MILLISECONDS) + ); + } + } + + private void sendResponseMessage(int uplinkMsgId, boolean success, String errorMsg) { + UplinkResponseMsg.Builder responseBuilder = UplinkResponseMsg.newBuilder() + .setUplinkMsgId(uplinkMsgId) + .setSuccess(success); + if (errorMsg != null) { + responseBuilder.setErrorMsg(errorMsg); + } + sendDownlinkMsg(ResponseMsg.newBuilder() + .setUplinkResponseMsg(responseBuilder.build()) + .build()); + } + + private void onDownlinkResponse(DownlinkResponseMsg msg) { + try { + if (msg.getSuccess()) { + sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId()); + log.debug("[{}][{}] Msg has been processed successfully! Msg Id: [{}], Msg: {}", tenantId, edge.getRoutingKey(), msg.getDownlinkMsgId(), msg); + } else { + log.error("[{}][{}] Msg processing failed! Msg Id: [{}], Error msg: {}", tenantId, edge.getRoutingKey(), msg.getDownlinkMsgId(), msg.getErrorMsg()); + } + if (sessionState.getPendingMsgsMap().isEmpty()) { + log.debug("[{}][{}] Pending msgs map is empty. Stopping current iteration", tenantId, edge.getRoutingKey()); + stopCurrentSendDownlinkMsgsTask(false); + } + } catch (Exception e) { + log.error("[{}][{}] Can't process downlink response message [{}]", tenantId, sessionId, msg, e); + } + } + + public void processHighPriorityEvents() { + try { + List highPriorityEvents = new ArrayList<>(); + EdgeEvent event; + while ((event = highPriorityQueue.poll()) != null) { + highPriorityEvents.add(event); + } + List downlinkMsgsPack = convertToDownlinkMsgsPack(highPriorityEvents); + sendDownlinkMsgsPack(downlinkMsgsPack).get(); + } catch (Exception e) { + log.error("[{}] Failed to process high priority events", sessionId, e); + } + } + + public ListenableFuture processEdgeEvents() throws Exception { + SettableFuture result = SettableFuture.create(); + log.trace("[{}][{}] starting processing edge events", tenantId, sessionId); + if (isConnected() && isSyncCompleted()) { + Pair startTsAndSeqId = getQueueStartTsAndSeqId().get(); + previousStartTs = startTsAndSeqId.getFirst(); + previousStartSeqId = startTsAndSeqId.getSecond(); + GeneralEdgeEventFetcher fetcher = new GeneralEdgeEventFetcher( + previousStartTs, + previousStartSeqId, + seqIdEnd, + false, + Integer.toUnsignedLong(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()), + ctx.getEdgeEventService()); + Futures.addCallback(startProcessingEdgeEvents(fetcher), new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Pair newStartTsAndSeqId) { + if (newStartTsAndSeqId != null) { + ListenableFuture> updateFuture = updateQueueStartTsAndSeqId(newStartTsAndSeqId); + Futures.addCallback(updateFuture, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable List list) { + log.debug("[{}][{}] queue offset was updated [{}]", tenantId, sessionId, newStartTsAndSeqId); + if (fetcher.isSeqIdNewCycleStarted()) { + seqIdEnd = fetcher.getSeqIdEnd(); + boolean newEventsAvailable = isNewEdgeEventsAvailable(); + result.set(newEventsAvailable); + } else { + seqIdEnd = null; + boolean newEventsAvailable = isSeqIdStartedNewCycle(); + if (!newEventsAvailable) { + newEventsAvailable = isNewEdgeEventsAvailable(); + } + result.set(newEventsAvailable); + } + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to update queue offset [{}]", tenantId, sessionId, newStartTsAndSeqId, t); + result.setException(t); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + log.trace("[{}][{}] newStartTsAndSeqId is null. Skipping iteration without db update", tenantId, sessionId); + result.set(Boolean.FALSE); + } + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to process events", tenantId, sessionId, t); + result.setException(t); + } + }, ctx.getGrpcCallbackExecutorService()); + } else { + log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId); + result.set(null); + } + return result; + } + + protected List convertToDownlinkMsgsPack(List edgeEvents) { + List result = new ArrayList<>(); + for (EdgeEvent edgeEvent : edgeEvents) { + log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, sessionId, edgeEvent); + DownlinkMsg downlinkMsg = null; + try { + switch (edgeEvent.getAction()) { + case UPDATED, ADDED, DELETED, ASSIGNED_TO_EDGE, UNASSIGNED_FROM_EDGE, ALARM_ACK, ALARM_CLEAR, + ALARM_DELETE, CREDENTIALS_UPDATED, RELATION_ADD_OR_UPDATE, RELATION_DELETED, RPC_CALL, + ASSIGNED_TO_CUSTOMER, UNASSIGNED_FROM_CUSTOMER, ADDED_COMMENT, UPDATED_COMMENT, DELETED_COMMENT -> { + downlinkMsg = convertEntityEventToDownlink(edgeEvent); + if (downlinkMsg != null && downlinkMsg.getWidgetTypeUpdateMsgCount() > 0) { + log.trace("[{}][{}] widgetTypeUpdateMsg message processed, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsgId()); + } else { + log.trace("[{}][{}] entity message processed [{}]", tenantId, sessionId, downlinkMsg); + } + } + case ATTRIBUTES_UPDATED, POST_ATTRIBUTES, ATTRIBUTES_DELETED, TIMESERIES_UPDATED -> + downlinkMsg = ctx.getTelemetryProcessor().convertTelemetryEventToDownlink(edge, edgeEvent); + default -> log.warn("[{}][{}] Unsupported action type [{}]", tenantId, sessionId, edgeEvent.getAction()); + } + } catch (Exception e) { + log.error("[{}][{}] Exception during converting edge event to downlink msg", tenantId, sessionId, e); + } + if (downlinkMsg != null) { + result.add(downlinkMsg); + } + } + return result; + } + + private ListenableFuture> getQueueStartTsAndSeqId() { + ListenableFuture> future = + ctx.getAttributesService().find(edge.getTenantId(), edge.getId(), AttributeScope.SERVER_SCOPE, Arrays.asList(QUEUE_START_TS_ATTR_KEY, QUEUE_START_SEQ_ID_ATTR_KEY)); + return Futures.transform(future, attributeKvEntries -> { + long startTs = 0L; + long startSeqId = 0L; + for (AttributeKvEntry attributeKvEntry : attributeKvEntries) { + if (QUEUE_START_TS_ATTR_KEY.equals(attributeKvEntry.getKey())) { + startTs = attributeKvEntry.getLongValue().isPresent() ? attributeKvEntry.getLongValue().get() : 0L; + } + if (QUEUE_START_SEQ_ID_ATTR_KEY.equals(attributeKvEntry.getKey())) { + startSeqId = attributeKvEntry.getLongValue().isPresent() ? attributeKvEntry.getLongValue().get() : 0L; + } + } + if (startSeqId == 0L) { + startSeqId = findStartSeqIdFromOldestEventIfAny(); + } + return Pair.of(startTs, startSeqId); + }, ctx.getGrpcCallbackExecutorService()); + } + + private boolean isSeqIdStartedNewCycle() { + try { + TimePageLink pageLink = new TimePageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount(), 0, null, null, newStartTs, System.currentTimeMillis()); + PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), 0L, previousStartSeqId == 0 ? null : previousStartSeqId - 1, pageLink); + return !edgeEvents.getData().isEmpty(); + } catch (Exception e) { + log.error("[{}][{}][{}] Failed to execute isSeqIdStartedNewCycle", tenantId, edge.getId(), sessionId, e); + } + return false; + } + + private boolean isNewEdgeEventsAvailable() { + try { + TimePageLink pageLink = new TimePageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount(), 0, null, null, newStartTs, System.currentTimeMillis()); + PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), newStartSeqId, null, pageLink); + return !edgeEvents.getData().isEmpty() || !highPriorityQueue.isEmpty(); + } catch (Exception e) { + log.error("[{}][{}][{}] Failed to execute isNewEdgeEventsAvailable", tenantId, edge.getId(), sessionId, e); + } + return false; + } + + private long findStartSeqIdFromOldestEventIfAny() { + long startSeqId = 0L; + try { + TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("createdTime"), null, null); + PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), null, null, pageLink); + if (!edgeEvents.getData().isEmpty()) { + startSeqId = edgeEvents.getData().get(0).getSeqId() - 1; + } + } catch (Exception e) { + log.error("[{}][{}][{}] Failed to execute findStartSeqIdFromOldestEventIfAny", tenantId, edge.getId(), sessionId, e); + } + return startSeqId; + } + + private ListenableFuture> updateQueueStartTsAndSeqId(Pair pair) { + newStartTs = pair.getFirst(); + newStartSeqId = pair.getSecond(); + log.trace("[{}] updateQueueStartTsAndSeqId [{}][{}][{}]", sessionId, edge.getId(), newStartTs, newStartSeqId); + List attributes = Arrays.asList( + new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_TS_ATTR_KEY, newStartTs), System.currentTimeMillis()), + new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_SEQ_ID_ATTR_KEY, newStartSeqId), System.currentTimeMillis())); + return ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), AttributeScope.SERVER_SCOPE, attributes); + } + + protected ListenableFuture> startProcessingEdgeEvents(EdgeEventFetcher fetcher) { + SettableFuture> result = SettableFuture.create(); + PageLink pageLink = fetcher.getPageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()); + processEdgeEvents(fetcher, pageLink, result); + return result; + } + + private void markSyncCompletedSendEdgeEventUpdate() { + syncCompleted = true; + ctx.getClusterService().onEdgeEventUpdate(new EdgeEventUpdateMsg(edge.getTenantId(), edge.getId())); + } + + private void stopCurrentSendDownlinkMsgsTask(Boolean isInterrupted) { + if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) { + sessionState.getSendDownlinkMsgsFuture().set(isInterrupted); + } + if (sessionState.getScheduledSendDownlinkTask() != null) { + sessionState.getScheduledSendDownlinkTask().cancel(true); + } + } + + private void sendDownlinkMsg(ResponseMsg downlinkMsg) { + if (downlinkMsg.getDownlinkMsg().getWidgetTypeUpdateMsgCount() > 0) { + log.trace("[{}][{}] Sending downlink widgetTypeUpdateMsg, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); + } else { + log.trace("[{}][{}] Sending downlink msg [{}]", tenantId, sessionId, downlinkMsg); + } + if (isConnected()) { + downlinkMsgLock.lock(); + try { + outputStream.onNext(downlinkMsg); + } catch (Exception e) { + log.error("[{}][{}] Failed to send downlink message [{}]", tenantId, sessionId, downlinkMsg, e); + connected = false; + sessionCloseListener.accept(edge, sessionId); + } finally { + downlinkMsgLock.unlock(); + } + log.trace("[{}][{}] Response msg successfully sent. downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); + } + } + + protected DownlinkMsg convertEntityEventToDownlink(EdgeEvent edgeEvent) { + log.trace("[{}] Executing convertEntityEventToDownlink, edgeEvent [{}], action [{}]", edgeEvent.getTenantId(), edgeEvent, edgeEvent.getAction()); + return switch (edgeEvent.getType()) { + case EDGE -> ctx.getEdgeProcessor().convertEdgeEventToDownlink(edgeEvent); + case DEVICE -> ctx.getDeviceProcessor().convertDeviceEventToDownlink(edgeEvent, edgeVersion); + case DEVICE_PROFILE -> ctx.getDeviceProfileProcessor().convertDeviceProfileEventToDownlink(edgeEvent, edgeVersion); + case ASSET_PROFILE -> ctx.getAssetProfileProcessor().convertAssetProfileEventToDownlink(edgeEvent, edgeVersion); + case ASSET -> ctx.getAssetProcessor().convertAssetEventToDownlink(edgeEvent, edgeVersion); + case ENTITY_VIEW -> ctx.getEntityViewProcessor().convertEntityViewEventToDownlink(edgeEvent, edgeVersion); + case DASHBOARD -> ctx.getDashboardProcessor().convertDashboardEventToDownlink(edgeEvent, edgeVersion); + case CUSTOMER -> ctx.getCustomerProcessor().convertCustomerEventToDownlink(edgeEvent, edgeVersion); + case RULE_CHAIN -> ctx.getRuleChainProcessor().convertRuleChainEventToDownlink(edgeEvent, edgeVersion); + case RULE_CHAIN_METADATA -> ctx.getRuleChainProcessor().convertRuleChainMetadataEventToDownlink(edgeEvent, edgeVersion); + case ALARM -> ctx.getAlarmProcessor().convertAlarmEventToDownlink(edgeEvent, edgeVersion); + case ALARM_COMMENT -> ctx.getAlarmProcessor().convertAlarmCommentEventToDownlink(edgeEvent, edgeVersion); + case USER -> ctx.getUserProcessor().convertUserEventToDownlink(edgeEvent, edgeVersion); + case RELATION -> ctx.getRelationProcessor().convertRelationEventToDownlink(edgeEvent, edgeVersion); + case WIDGETS_BUNDLE -> ctx.getWidgetBundleProcessor().convertWidgetsBundleEventToDownlink(edgeEvent, edgeVersion); + case WIDGET_TYPE -> ctx.getWidgetTypeProcessor().convertWidgetTypeEventToDownlink(edgeEvent, edgeVersion); + case ADMIN_SETTINGS -> ctx.getAdminSettingsProcessor().convertAdminSettingsEventToDownlink(edgeEvent, edgeVersion); + case OTA_PACKAGE -> ctx.getOtaPackageProcessor().convertOtaPackageEventToDownlink(edgeEvent, edgeVersion); + case TB_RESOURCE -> ctx.getResourceProcessor().convertResourceEventToDownlink(edgeEvent, edgeVersion); + case QUEUE -> ctx.getQueueProcessor().convertQueueEventToDownlink(edgeEvent, edgeVersion); + case TENANT -> ctx.getTenantProcessor().convertTenantEventToDownlink(edgeEvent, edgeVersion); + case TENANT_PROFILE -> ctx.getTenantProfileProcessor().convertTenantProfileEventToDownlink(edgeEvent, edgeVersion); + case NOTIFICATION_RULE -> ctx.getNotificationEdgeProcessor().convertNotificationRuleToDownlink(edgeEvent); + case NOTIFICATION_TARGET -> ctx.getNotificationEdgeProcessor().convertNotificationTargetToDownlink(edgeEvent); + case NOTIFICATION_TEMPLATE -> ctx.getNotificationEdgeProcessor().convertNotificationTemplateToDownlink(edgeEvent); + case OAUTH2_CLIENT -> ctx.getOAuth2EdgeProcessor().convertOAuth2ClientEventToDownlink(edgeEvent, edgeVersion); + case DOMAIN -> ctx.getOAuth2EdgeProcessor().convertOAuth2DomainEventToDownlink(edgeEvent, edgeVersion); + default -> { + log.warn("[{}] Unsupported edge event type [{}]", edgeEvent.getTenantId(), edgeEvent); + yield null; + } + }; + } + + public void addEventToHighPriorityQueue(EdgeEvent edgeEvent) { + while (highPriorityQueue.size() > maxHighPriorityQueueSizePerSession) { + EdgeEvent oldestHighPriority = highPriorityQueue.poll(); + if (oldestHighPriority != null) { + log.warn("[{}][{}][{}] High priority queue is full. Removing oldest high priority event from queue {}", + tenantId, edge.getId(), sessionId, oldestHighPriority); + } + } + highPriorityQueue.add(edgeEvent); + } + + protected ListenableFuture> processUplinkMsg(UplinkMsg uplinkMsg) { + List> result = new ArrayList<>(); + try { + if (uplinkMsg.getDeviceProfileUpdateMsgCount() > 0) { + for (DeviceProfileUpdateMsg deviceProfileUpdateMsg : uplinkMsg.getDeviceProfileUpdateMsgList()) { + result.add(((DeviceProfileProcessor) ctx.getDeviceProfileEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processDeviceProfileMsgFromEdge(edge.getTenantId(), edge, deviceProfileUpdateMsg)); + } + } + if (uplinkMsg.getDeviceUpdateMsgCount() > 0) { + for (DeviceUpdateMsg deviceUpdateMsg : uplinkMsg.getDeviceUpdateMsgList()) { + result.add(((DeviceProcessor) ctx.getDeviceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processDeviceMsgFromEdge(edge.getTenantId(), edge, deviceUpdateMsg)); + } + } + if (uplinkMsg.getDeviceCredentialsUpdateMsgCount() > 0) { + for (DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg : uplinkMsg.getDeviceCredentialsUpdateMsgList()) { + result.add(((DeviceProcessor) ctx.getDeviceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processDeviceCredentialsMsgFromEdge(edge.getTenantId(), edge.getId(), deviceCredentialsUpdateMsg)); + } + } + if (uplinkMsg.getAssetProfileUpdateMsgCount() > 0) { + for (AssetProfileUpdateMsg assetProfileUpdateMsg : uplinkMsg.getAssetProfileUpdateMsgList()) { + result.add(((AssetProfileProcessor) ctx.getAssetProfileEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processAssetProfileMsgFromEdge(edge.getTenantId(), edge, assetProfileUpdateMsg)); + } + } + if (uplinkMsg.getAssetUpdateMsgCount() > 0) { + for (AssetUpdateMsg assetUpdateMsg : uplinkMsg.getAssetUpdateMsgList()) { + result.add(((AssetProcessor) ctx.getAssetEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processAssetMsgFromEdge(edge.getTenantId(), edge, assetUpdateMsg)); + } + } + if (uplinkMsg.getEntityViewUpdateMsgCount() > 0) { + for (EntityViewUpdateMsg entityViewUpdateMsg : uplinkMsg.getEntityViewUpdateMsgList()) { + result.add(((EntityViewProcessor) ctx.getEntityViewProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processEntityViewMsgFromEdge(edge.getTenantId(), edge, entityViewUpdateMsg)); + } + } + if (uplinkMsg.getEntityDataCount() > 0) { + for (EntityDataProto entityData : uplinkMsg.getEntityDataList()) { + result.addAll(ctx.getTelemetryProcessor().processTelemetryMsg(edge.getTenantId(), entityData)); + } + } + if (uplinkMsg.getAlarmUpdateMsgCount() > 0) { + for (AlarmUpdateMsg alarmUpdateMsg : uplinkMsg.getAlarmUpdateMsgList()) { + result.add(((AlarmProcessor) ctx.getAlarmEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processAlarmMsgFromEdge(edge.getTenantId(), edge.getId(), alarmUpdateMsg)); + } + } + if (uplinkMsg.getAlarmCommentUpdateMsgCount() > 0) { + for (AlarmCommentUpdateMsg alarmCommentUpdateMsg : uplinkMsg.getAlarmCommentUpdateMsgList()) { + result.add(((AlarmProcessor) ctx.getAlarmEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processAlarmCommentMsgFromEdge(edge.getTenantId(), edge.getId(), alarmCommentUpdateMsg)); + } + } + if (uplinkMsg.getRelationUpdateMsgCount() > 0) { + for (RelationUpdateMsg relationUpdateMsg : uplinkMsg.getRelationUpdateMsgList()) { + result.add(((RelationProcessor) ctx.getRelationEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processRelationMsgFromEdge(edge.getTenantId(), edge, relationUpdateMsg)); + } + } + if (uplinkMsg.getDashboardUpdateMsgCount() > 0) { + for (DashboardUpdateMsg dashboardUpdateMsg : uplinkMsg.getDashboardUpdateMsgList()) { + result.add(((DashboardProcessor) ctx.getDashboardEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processDashboardMsgFromEdge(edge.getTenantId(), edge, dashboardUpdateMsg)); + } + } + if (uplinkMsg.getResourceUpdateMsgCount() > 0) { + for (ResourceUpdateMsg resourceUpdateMsg : uplinkMsg.getResourceUpdateMsgList()) { + result.add(((ResourceProcessor) ctx.getResourceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processResourceMsgFromEdge(edge.getTenantId(), edge, resourceUpdateMsg)); + } + } + if (uplinkMsg.getRuleChainMetadataRequestMsgCount() > 0) { + for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processRuleChainMetadataRequestMsg(edge.getTenantId(), edge, ruleChainMetadataRequestMsg)); + } + } + if (uplinkMsg.getAttributesRequestMsgCount() > 0) { + for (AttributesRequestMsg attributesRequestMsg : uplinkMsg.getAttributesRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processAttributesRequestMsg(edge.getTenantId(), edge, attributesRequestMsg)); + } + } + if (uplinkMsg.getRelationRequestMsgCount() > 0) { + for (RelationRequestMsg relationRequestMsg : uplinkMsg.getRelationRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processRelationRequestMsg(edge.getTenantId(), edge, relationRequestMsg)); + } + } + if (uplinkMsg.getUserCredentialsRequestMsgCount() > 0) { + for (UserCredentialsRequestMsg userCredentialsRequestMsg : uplinkMsg.getUserCredentialsRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processUserCredentialsRequestMsg(edge.getTenantId(), edge, userCredentialsRequestMsg)); + } + } + if (uplinkMsg.getDeviceCredentialsRequestMsgCount() > 0) { + for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg : uplinkMsg.getDeviceCredentialsRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processDeviceCredentialsRequestMsg(edge.getTenantId(), edge, deviceCredentialsRequestMsg)); + } + } + if (uplinkMsg.getDeviceRpcCallMsgCount() > 0) { + for (DeviceRpcCallMsg deviceRpcCallMsg : uplinkMsg.getDeviceRpcCallMsgList()) { + result.add(((DeviceProcessor) ctx.getDeviceEdgeProcessorFactory().getProcessorByEdgeVersion(edgeVersion)) + .processDeviceRpcCallFromEdge(edge.getTenantId(), edge, deviceRpcCallMsg)); + } + } + if (uplinkMsg.getWidgetBundleTypesRequestMsgCount() > 0) { + for (WidgetBundleTypesRequestMsg widgetBundleTypesRequestMsg : uplinkMsg.getWidgetBundleTypesRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processWidgetBundleTypesRequestMsg(edge.getTenantId(), edge, widgetBundleTypesRequestMsg)); + } + } + if (uplinkMsg.getEntityViewsRequestMsgCount() > 0) { + for (EntityViewsRequestMsg entityViewRequestMsg : uplinkMsg.getEntityViewsRequestMsgList()) { + result.add(ctx.getEdgeRequestsService().processEntityViewsRequestMsg(edge.getTenantId(), edge, entityViewRequestMsg)); + } + } + } catch (Exception e) { + String failureMsg = String.format("Can't process uplink msg [%s] from edge", uplinkMsg); + log.error("[{}][{}] Can't process uplink msg [{}]", edge.getTenantId(), sessionId, uplinkMsg, e); + ctx.getNotificationRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(edge.getTenantId()).edgeId(edge.getId()) + .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(e.getMessage()).build()); + return Futures.immediateFailedFuture(e); + } + return Futures.allAsList(result); + } + + protected void destroy() { + // used for KafkaEdgeGrpcSession only + } + + protected void deleteTopic(EdgeId edgeId) { + // used for KafkaEdgeGrpcSession only + } + + @Override + public void close() { + log.debug("[{}][{}] Closing session", tenantId, sessionId); + connected = false; + try { + outputStream.onCompleted(); + } catch (Exception e) { + log.debug("[{}][{}] Failed to close output stream: {}", tenantId, sessionId, e.getMessage()); + } + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java index a65e6ab85f..6aaa3d72a4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java @@ -30,6 +30,10 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificat import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; +import org.thingsboard.server.queue.discovery.TopicService; +import org.thingsboard.server.queue.kafka.TbKafkaAdmin; +import org.thingsboard.server.queue.kafka.TbKafkaSettings; +import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.service.edge.EdgeContextComponent; @@ -42,22 +46,29 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.function.BiConsumer; @Slf4j -public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { +public class KafkaEdgeGrpcSession extends EdgeGrpcSession { + private final TopicService topicService; private final TbCoreQueueFactory tbCoreQueueFactory; + private final TbKafkaSettings kafkaSettings; + private final TbKafkaTopicConfigs kafkaTopicConfigs; + private volatile boolean isHighPriorityProcessing; private QueueConsumerManager> consumer; private ExecutorService consumerExecutor; - public KafkaEdgeGrpcSession(EdgeContextComponent ctx, TbCoreQueueFactory tbCoreQueueFactory, StreamObserver outputStream, - BiConsumer sessionOpenListener, - BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, - int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { + public KafkaEdgeGrpcSession(EdgeContextComponent ctx, TopicService topicService, TbCoreQueueFactory tbCoreQueueFactory, + TbKafkaSettings kafkaSettings, TbKafkaTopicConfigs kafkaTopicConfigs, StreamObserver outputStream, + BiConsumer sessionOpenListener, BiConsumer sessionCloseListener, + ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); + this.topicService = topicService; this.tbCoreQueueFactory = tbCoreQueueFactory; + this.kafkaSettings = kafkaSettings; + this.kafkaTopicConfigs = kafkaTopicConfigs; } private void processMsgs(List> msgs, TbQueueConsumer> consumer) { @@ -125,4 +136,11 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { consumerExecutor.shutdown(); } + @Override + public void deleteTopic(EdgeId edgeId) { + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); + kafkaAdmin.deleteTopic(topic); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java index cddd895c02..12516568a8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java @@ -29,19 +29,16 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.function.BiConsumer; @Slf4j -public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession { +public class PostgresEdgeGrpcSession extends EdgeGrpcSession { PostgresEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, - BiConsumer sessionOpenListener, + BiConsumer sessionOpenListener, BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); initInputStream(); } - @Override - public void destroy() {} - @Override public ListenableFuture migrateEdgeEvents() { return Futures.immediateFuture(Boolean.FALSE); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index 38cebbb509..bce599c0e3 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java @@ -35,7 +35,6 @@ import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; @@ -55,12 +54,7 @@ import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.tenant.TenantService; -import org.thingsboard.server.queue.discovery.TopicService; -import org.thingsboard.server.queue.kafka.TbKafkaAdmin; -import org.thingsboard.server.queue.kafka.TbKafkaSettings; -import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs; -import java.util.Optional; import java.util.Set; @Slf4j @@ -68,12 +62,8 @@ import java.util.Set; @RequiredArgsConstructor public class EntityStateSourcingListener { - private final TopicService topicService; - private final TbClusterService tbClusterService; private final TenantService tenantService; - - private final Optional kafkaSettings; - private final Optional kafkaTopicConfigs; + private final TbClusterService tbClusterService; @PostConstruct public void init() { @@ -147,7 +137,7 @@ public class EntityStateSourcingListener { log.debug("[{}][{}][{}] Handling entity deletion event: {}", tenantId, entityType, entityId, event); switch (entityType) { - case ASSET, ASSET_PROFILE, ENTITY_VIEW, CUSTOMER, NOTIFICATION_RULE -> { + case ASSET, ASSET_PROFILE, EDGE, ENTITY_VIEW, CUSTOMER, NOTIFICATION_RULE -> { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); } case NOTIFICATION_REQUEST -> { @@ -187,10 +177,6 @@ public class EntityStateSourcingListener { TbResourceInfo tbResource = (TbResourceInfo) event.getEntity(); tbClusterService.onResourceDeleted(tbResource, null); } - case EDGE -> { - onEdgeDelete(tenantId, (EdgeId) event.getEntityId()); - tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); - } default -> {} } } @@ -261,14 +247,6 @@ public class EntityStateSourcingListener { } } - private void onEdgeDelete(TenantId tenantId, EdgeId edgeId) { - if (kafkaSettings.isPresent() && kafkaTopicConfigs.isPresent()) { - String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); - TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings.get(), kafkaTopicConfigs.get().getEdgeEventConfigs()); - kafkaAdmin.deleteTopic(topic); - } - } - private void pushAssignedFromNotification(Tenant currentTenant, TenantId newTenantId, Device assignedDevice) { String data = JacksonUtil.toString(JacksonUtil.valueToTree(assignedDevice)); if (data != null) { diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java index dfd6226ed3..1bd0e30c77 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java @@ -65,20 +65,16 @@ public class KafkaEdgeTopicsCleanUpService { @Value("${sql.ttl.edge_events.edge_events_ttl:2628000}") private long ttlSeconds; - private final ExecutorService executorService = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("kafka-edge-topic-cleanup")); - @Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.edge_events.execution_interval_ms})}", fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}") public void cleanUp() { - executorService.submit(() -> { - PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); - for (TenantId tenantId : tenants) { - try { - cleanUp(tenantId); - } catch (Exception e) { - log.warn("Failed to drop kafka topics for tenant {}", tenantId, e); - } + PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); + for (TenantId tenantId : tenants) { + try { + cleanUp(tenantId); + } catch (Exception e) { + log.warn("Failed to drop kafka topics for tenant {}", tenantId, e); } - }); + } } private void cleanUp(TenantId tenantId) throws Exception { @@ -91,12 +87,12 @@ public class KafkaEdgeTopicsCleanUpService { long ttlMillis = TimeUnit.SECONDS.toChronoUnit().getDuration().multipliedBy(ttlSeconds).toMillis(); for (EdgeId edgeId : edgeIds) { + TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); attributesService.find(tenantId, edgeId, AttributeScope.SERVER_SCOPE, DefaultDeviceStateService.LAST_CONNECT_TIME).get() .flatMap(AttributeKvEntry::getLongValue) .filter(lastConnectTime -> isTopicExpired(lastConnectTime, ttlMillis, currentTimeMillis)) .ifPresent(lastConnectTime -> { String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); - TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); if (kafkaAdmin.isTopicEmpty(topic)) { kafkaAdmin.deleteTopic(topic); log.info("Removed outdated topic for tenant {} and edge with id {} older than {}", diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java index 2540cdeafa..7960cbcf32 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java @@ -83,7 +83,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider @Override public TbQueueProducer> getTbEdgeMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used Transport!"); + throw new RuntimeException("Not Implemented! Should not be used by Transport!"); } @Override @@ -93,7 +93,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider @Override public TbQueueProducer> getTbEdgeEventsMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used Transport!"); + throw new RuntimeException("Not Implemented! Should not be used by Transport!"); } @Override