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 index 9aa21dd265..0baa8a90a4 100644 --- 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 @@ -115,16 +115,16 @@ public abstract class AbstractEdgeGrpcSession highPriorityQueue = new ConcurrentLinkedQueue<>(); - private UUID sessionId; + protected UUID sessionId; private BiConsumer sessionOpenListener; private BiConsumer sessionCloseListener; private final EdgeSessionState sessionState = new EdgeSessionState(); private final ReentrantLock downlinkMsgLock = new ReentrantLock(); - private EdgeContextComponent ctx; - private Edge edge; - private TenantId tenantId; + protected EdgeContextComponent ctx; + protected Edge edge; + protected TenantId tenantId; private Long newStartTs; private Long previousStartTs; @@ -162,7 +162,7 @@ public abstract class AbstractEdgeGrpcSession() { + inputStream = new StreamObserver<>() { @Override public void onNext(RequestMsg requestMsg) { if (!connected && requestMsg.getMsgType().equals(RequestMsgType.CONNECT_RPC_MESSAGE)) { @@ -233,9 +233,9 @@ public abstract class AbstractEdgeGrpcSession> future = startProcessingEdgeEvents(next); Futures.addCallback(future, new FutureCallback<>() { @Override @@ -270,7 +270,7 @@ public abstract class AbstractEdgeGrpcSession pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); if (isConnected() && !pageData.getData().isEmpty()) { - log.trace("[{}][{}][{}] event(s) are going to be processed.", this.tenantId, this.sessionId, pageData.getData().size()); + log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, sessionId, pageData.getData().size()); List downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); + for (DownlinkMsg downlinkMsg : downlinkMsgsPack) { + if (downlinkMsg.getEntityDataCount() > 0) { + System.out.println("downlink = " + downlinkMsg); + } + } Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { @Override public void onSuccess(@Nullable Boolean isInterrupted) { @@ -329,25 +334,25 @@ public abstract class AbstractEdgeGrpcSession optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); if (optional.isPresent()) { edge = optional.get(); - this.tenantId = edge.getTenantId(); + tenantId = edge.getTenantId(); try { if (edge.getSecret().equals(request.getEdgeSecret())) { sessionOpenListener.accept(edge.getId(), (T) this); - this.edgeVersion = request.getEdgeVersion(); + edgeVersion = request.getEdgeVersion(); processSaveEdgeVersionAsAttribute(request.getEdgeVersion().name()); return ConnectResponseMsg.newBuilder() .setResponseCode(ConnectResponseCode.ACCEPTED) @@ -383,11 +388,11 @@ public abstract class AbstractEdgeGrpcSession this.clientMaxInboundMessageSize) { + 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.", this.clientMaxInboundMessageSize); + "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(), this.clientMaxInboundMessageSize); - log.error("[{}][{}][{}] {} Message {}", this.tenantId, edge.getId(), this.sessionId, message, downlinkMsg); + "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()); @@ -478,7 +483,7 @@ public abstract class AbstractEdgeGrpcSession downlinkMsgsPack = convertToDownlinkMsgsPack(highPriorityEvents); sendDownlinkMsgsPack(downlinkMsgsPack).get(); } catch (Exception e) { - log.error("[{}] Failed to process high priority events", this.sessionId, e); + log.error("[{}] Failed to process high priority events", sessionId, e); } } protected ListenableFuture processEdgeEvents() throws Exception { SettableFuture result = SettableFuture.create(); - log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId); + log.trace("[{}][{}] starting processing edge events", tenantId, sessionId); if (isConnected() && isSyncCompleted()) { Pair startTsAndSeqId = getQueueStartTsAndSeqId().get(); - this.previousStartTs = startTsAndSeqId.getFirst(); - this.previousStartSeqId = startTsAndSeqId.getSecond(); + previousStartTs = startTsAndSeqId.getFirst(); + previousStartSeqId = startTsAndSeqId.getSecond(); GeneralEdgeEventFetcher fetcher = new GeneralEdgeEventFetcher( - this.previousStartTs, - this.previousStartSeqId, - this.seqIdEnd, + previousStartTs, + previousStartSeqId, + seqIdEnd, false, Integer.toUnsignedLong(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()), ctx.getEdgeEventService()); @@ -593,7 +598,7 @@ public abstract class AbstractEdgeGrpcSession convertToDownlinkMsgsPack(List edgeEvents) { List result = new ArrayList<>(); for (EdgeEvent edgeEvent : edgeEvents) { - log.trace("[{}][{}] converting edge event to downlink msg [{}]", this.tenantId, this.sessionId, edgeEvent); + log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, sessionId, edgeEvent); DownlinkMsg downlinkMsg = null; try { switch (edgeEvent.getAction()) { @@ -624,18 +629,18 @@ public abstract class AbstractEdgeGrpcSession { downlinkMsg = convertEntityEventToDownlink(edgeEvent); if (downlinkMsg != null && downlinkMsg.getWidgetTypeUpdateMsgCount() > 0) { - log.trace("[{}][{}] widgetTypeUpdateMsg message processed, downlinkMsgId = {}", this.tenantId, this.sessionId, downlinkMsg.getDownlinkMsgId()); + log.trace("[{}][{}] widgetTypeUpdateMsg message processed, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsgId()); } else { - log.trace("[{}][{}] entity message processed [{}]", this.tenantId, this.sessionId, downlinkMsg); + 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 [{}]", this.tenantId, this.sessionId, edgeEvent.getAction()); + log.warn("[{}][{}] Unsupported action type [{}]", tenantId, sessionId, edgeEvent.getAction()); } } catch (Exception e) { - log.error("[{}][{}] Exception during converting edge event to downlink msg", this.tenantId, this.sessionId, e); + log.error("[{}][{}] Exception during converting edge event to downlink msg", tenantId, sessionId, e); } if (downlinkMsg != null) { result.add(downlinkMsg); @@ -667,22 +672,22 @@ public abstract class AbstractEdgeGrpcSession edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), 0L, this.previousStartSeqId == 0 ? null : this.previousStartSeqId - 1, pageLink); + 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", this.tenantId, edge.getId(), sessionId, 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, this.newStartTs, System.currentTimeMillis()); - PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), this.newStartSeqId, null, pageLink); + 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", this.tenantId, edge.getId(), sessionId, e); + log.error("[{}][{}][{}] Failed to execute isNewEdgeEventsAvailable", tenantId, edge.getId(), sessionId, e); } return false; } @@ -696,18 +701,18 @@ public abstract class AbstractEdgeGrpcSession> updateQueueStartTsAndSeqId(Pair pair) { - this.newStartTs = pair.getFirst(); - this.newStartSeqId = pair.getSecond(); - log.trace("[{}] updateQueueStartTsAndSeqId [{}][{}][{}]", this.sessionId, edge.getId(), this.newStartTs, this.newStartSeqId); + 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, this.newStartTs), System.currentTimeMillis()), - new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_SEQ_ID_ATTR_KEY, this.newStartSeqId), System.currentTimeMillis())); + 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); } @@ -734,22 +739,22 @@ public abstract class AbstractEdgeGrpcSession 0) { - log.trace("[{}][{}] Sending downlink widgetTypeUpdateMsg, downlinkMsgId = {}", this.tenantId, this.sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); + log.trace("[{}][{}] Sending downlink widgetTypeUpdateMsg, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); } else { - log.trace("[{}][{}] Sending downlink msg [{}]", this.tenantId, this.sessionId, downlinkMsg); + 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 [{}]", this.tenantId, this.sessionId, downlinkMsg, 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 = {}", this.tenantId, this.sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); + log.trace("[{}][{}] Response msg successfully sent. downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsg().getDownlinkMsgId()); } } @@ -805,7 +810,7 @@ public abstract class AbstractEdgeGrpcSession session) { EdgeId edgeId = session.getEdge().getId(); - UUID tenantId = session.getEdge().getTenantId().getId(); + TenantId tenantId = session.getEdge().getTenantId(); if (sessions.containsKey(edgeId)) { ScheduledFuture edgeEventCheckTask = edgeEventProcessingExecutorService.schedule(() -> { try { @@ -423,8 +422,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (Boolean.TRUE.equals(newEventsAdded)) { sessionNewEvents.put(edgeId, true); } - if (isKafkaSupported) { + if (session instanceof KafkaEdgeGrpcSession kafkaEdgeGrpcSession && newEventsAdded != null) { edgeEventsProcessed.put(edgeId, true); + initializeKafkaConsumer(kafkaEdgeGrpcSession, tenantId, edgeId); + cancelScheduleEdgeEventsCheck(edgeId); } else { scheduleEdgeEventsCheck(session); } @@ -437,7 +438,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } }, ctx.getGrpcCallbackExecutorService()); } else { - if (!isKafkaSupported) { + if (Boolean.TRUE.equals(edgeEventsProcessed.get(edgeId)) && session instanceof KafkaEdgeGrpcSession kafkaEdgeGrpcSession) { + initializeKafkaConsumer(kafkaEdgeGrpcSession, tenantId, edgeId); + cancelScheduleEdgeEventsCheck(edgeId); + } else { scheduleEdgeEventsCheck(session); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java index 77760524a5..1126214e60 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java @@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.util.ProtoUtils; +import org.thingsboard.server.dao.edge.EdgeEventDao; import org.thingsboard.server.dao.edge.EdgeEventService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; @@ -54,6 +55,7 @@ public class KafkaEdgeEventService implements EdgeEventService { private final TbQueueProducerProvider producerProvider; @Lazy private final TopicService topicService; + private final EdgeEventDao edgeEventDao; @Override public ListenableFuture saveAsync(EdgeEvent edgeEvent) { @@ -74,11 +76,14 @@ public class KafkaEdgeEventService implements EdgeEventService { @Override public PageData findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) { - return null; + // To support fetching edge events on connect from postgres if there are any: + return edgeEventDao.findEdgeEvents(tenantId.getId(), edgeId, seqIdStart, seqIdEnd, pageLink); } @Override public void cleanupEvents(long ttl) { + // To delete deprecated events by ttl + edgeEventDao.cleanupEvents(ttl); } } 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 b86aa0ecf2..920047898b 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 @@ -71,7 +71,7 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession> msgs, TbQueueConsumer> consumer) { - log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId); + log.trace("[{}][{}] starting processing edge events", tenantId, sessionId); if (isConnected() && isSyncCompleted()) { if (!highPriorityQueue.isEmpty()) {