diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index 3653eb3a40..9317c51a64 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -20,10 +20,6 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import jakarta.annotation.PostConstruct; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; -import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionalEventListener; import org.thingsboard.common.util.JacksonUtil; @@ -41,7 +37,6 @@ import org.thingsboard.server.common.data.domain.Domain; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; -import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.relation.EntityRelation; @@ -55,10 +50,6 @@ import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.RelationActionEvent; 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; /** * This event listener does not support async event processing because relay on ThreadLocal @@ -79,22 +70,11 @@ import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs; @RequiredArgsConstructor public class EdgeEventSourcingListener { - private final TopicService topicService; private final TbClusterService tbClusterService; private final TenantService tenantService; private final EdgeSynchronizationManager edgeSynchronizationManager; - @Autowired(required = false) - @Lazy - private TbKafkaSettings kafkaSettings; - @Autowired(required = false) - @Lazy - private TbKafkaTopicConfigs kafkaTopicConfigs; - - @Value("#{'${queue.type:null}' == 'kafka'}") - private boolean isKafkaSupported; - @PostConstruct public void init() { log.debug("EdgeEventSourcingListener initiated"); @@ -127,11 +107,7 @@ public class EdgeEventSourcingListener { return; } try { - if (EntityType.TENANT.equals(entityType)) { - return; - } - if (EntityType.EDGE.equals(entityType)) { - handleEdgeEntityDeletion((EdgeId) event.getEntityId(), tenantId); + if (EntityType.TENANT.equals(entityType) || EntityType.EDGE.equals(entityType)) { return; } log.trace("[{}] DeleteEntityEvent called: {}", tenantId, event); @@ -145,14 +121,6 @@ public class EdgeEventSourcingListener { } } - private void handleEdgeEntityDeletion(EdgeId edgeId, TenantId tenantId) { - if (isKafkaSupported) { - String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); - TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); - kafkaAdmin.deleteTopic(topic); - } - } - private EdgeEventActionType getEdgeEventActionTypeForEntityEvent(Object entity) { if (entity instanceof AlarmComment) { return EdgeEventActionType.DELETED_COMMENT; 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 15ee44a66e..2f06b4d6e4 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 @@ -539,7 +539,8 @@ public abstract class AbstractEdgeGrpcSession highPriorityEvents = new ArrayList<>(); EdgeEvent event; @@ -553,7 +554,8 @@ public abstract class AbstractEdgeGrpcSession processEdgeEvents() throws Exception { + @Override + public ListenableFuture processEdgeEvents() throws Exception { SettableFuture result = SettableFuture.create(); log.trace("[{}][{}] starting processing edge events", tenantId, sessionId); if (isConnected() && isSyncCompleted()) { 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 a82ba3efb2..dbb2166d69 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 @@ -57,9 +57,6 @@ import org.thingsboard.server.common.msg.edge.ToEdgeSyncRequest; 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.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; -import org.thingsboard.server.queue.TbQueueConsumer; -import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.TbCoreComponent; @@ -95,7 +92,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private final ConcurrentMap> sessionEdgeEventChecks = new ConcurrentHashMap<>(); private final ConcurrentMap> localSyncEdgeRequests = new ConcurrentHashMap<>(); private final ConcurrentMap edgeEventsProcessed = new ConcurrentHashMap<>(); - private final ConcurrentMap kafkaConsumerInit = new ConcurrentHashMap<>(); @Value("${edges.rpc.port}") private int rpcPort; @@ -220,7 +216,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private AbstractEdgeGrpcSession createEdgeGrpcSession(StreamObserver outputStream) { return isKafkaSupported - ? new KafkaEdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, + ? new KafkaEdgeGrpcSession(ctx, tbCoreQueueFactory, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession) : new PostgresEdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); @@ -333,26 +329,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i pushRuleEngineMessage(tenantId, edge, lastConnectTs, TbMsgType.CONNECT_EVENT); cancelScheduleEdgeEventsCheck(edgeId); edgeEventsProcessed.putIfAbsent(edgeId, Boolean.FALSE); - if (edgeGrpcSession instanceof KafkaEdgeGrpcSession session) { - Boolean isChecked = edgeEventsProcessed.get(edgeId); - if (Boolean.FALSE.equals(isChecked)) { - scheduleEdgeEventsCheck(session); - } else { - initializeKafkaConsumer(session, tenantId, edgeId); - } - } - if (edgeGrpcSession instanceof PostgresEdgeGrpcSession) { - scheduleEdgeEventsCheck(edgeGrpcSession); - } - } - - private void initializeKafkaConsumer(KafkaEdgeGrpcSession kafkaEdgeGrpcSession, TenantId tenantId, EdgeId edgeId) { - TbQueueConsumer> consumer = tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edgeId); - if (!kafkaConsumerInit.getOrDefault(edgeId, Boolean.FALSE)) { - kafkaEdgeGrpcSession.initConsumer(() -> consumer, schedulerPoolSize); - kafkaEdgeGrpcSession.startConsumers(); - kafkaConsumerInit.put(edgeId, Boolean.TRUE); - } + scheduleEdgeEventsCheck(edgeGrpcSession); } private void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId, String requestServiceId) { @@ -419,22 +396,18 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); newEventLock.lock(); try { - if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId)) && Boolean.FALSE.equals(edgeEventsProcessed.get(edgeId))) { + if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { log.trace("[{}][{}] Set session new events flag to false", tenantId, edgeId.getId()); sessionNewEvents.put(edgeId, false); + processEdgeEventMigrationIfNeeded(session, edgeId); + session.processHighPriorityEvents(); Futures.addCallback(session.processEdgeEvents(), new FutureCallback<>() { @Override public void onSuccess(Boolean newEventsAdded) { if (Boolean.TRUE.equals(newEventsAdded)) { sessionNewEvents.put(edgeId, true); } - if (session instanceof KafkaEdgeGrpcSession kafkaEdgeGrpcSession && newEventsAdded != null) { - edgeEventsProcessed.put(edgeId, true); - initializeKafkaConsumer(kafkaEdgeGrpcSession, tenantId, edgeId); - cancelScheduleEdgeEventsCheck(edgeId); - } else { - scheduleEdgeEventsCheck(session); - } + scheduleEdgeEventsCheck(session); } @Override @@ -444,12 +417,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } }, ctx.getGrpcCallbackExecutorService()); } else { - if (Boolean.TRUE.equals(edgeEventsProcessed.get(edgeId)) && session instanceof KafkaEdgeGrpcSession kafkaEdgeGrpcSession) { - initializeKafkaConsumer(kafkaEdgeGrpcSession, tenantId, edgeId); - cancelScheduleEdgeEventsCheck(edgeId); - } else { - scheduleEdgeEventsCheck(session); - } + scheduleEdgeEventsCheck(session); } } finally { newEventLock.unlock(); @@ -466,6 +434,21 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } + private void processEdgeEventMigrationIfNeeded(AbstractEdgeGrpcSession session, EdgeId edgeId) throws Exception { + boolean isMigrationProcessed = edgeEventsProcessed.getOrDefault(edgeId, Boolean.FALSE); + if (!isMigrationProcessed) { + Boolean migrated = session.migrateEdgeEvents(false).get(); + if (Boolean.TRUE.equals(migrated)) { + sessionNewEvents.put(edgeId, true); + scheduleEdgeEventsCheck(session); + } else if (Boolean.FALSE.equals(migrated)) { + edgeEventsProcessed.put(edgeId, true); + } else { + scheduleEdgeEventsCheck(session); + } + } + } + private void cancelScheduleEdgeEventsCheck(EdgeId edgeId) { log.trace("[{}] cancelling edge event check for edge", edgeId); if (sessionEdgeEventChecks.containsKey(edgeId)) { @@ -492,7 +475,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } if (isKafkaSupported) { ((KafkaEdgeGrpcSession) toRemove).stopConsumer(); - kafkaConsumerInit.remove(edgeId); } TenantId tenantId = toRemove.getEdge().getTenantId(); save(tenantId, edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); 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 6a72b40b80..3db3a72288 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,6 +15,7 @@ */ package org.thingsboard.server.service.edge.rpc; +import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.edge.Edge; public interface EdgeGrpcSession { @@ -25,4 +26,10 @@ public interface EdgeGrpcSession { boolean isConnected(); + ListenableFuture migrateEdgeEvents(boolean isMigrationProcessed) throws Exception; + + ListenableFuture processEdgeEvents() throws Exception; + + void processHighPriorityEvents(); + } 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 e16c89c23e..4fabfd5945 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 @@ -15,6 +15,8 @@ */ package org.thingsboard.server.service.edge.rpc; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import io.grpc.stub.StreamObserver; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; @@ -29,6 +31,7 @@ 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.provider.TbCoreQueueFactory; import org.thingsboard.server.service.edge.EdgeContextComponent; import java.util.ArrayList; @@ -37,86 +40,50 @@ import java.util.UUID; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import java.util.function.BiConsumer; -import java.util.function.Supplier; @Slf4j public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { - private ExecutorService consumerExecutor; - private ScheduledExecutorService highPriorityExecutorService; + private final TbQueueConsumer> edgeEventsConsumer; + + private volatile boolean isHighPriorityProcessing; + private volatile boolean isConsumerInit; private QueueConsumerManager> consumer; - public KafkaEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, + private ExecutorService consumerExecutor; + + public KafkaEdgeGrpcSession(EdgeContextComponent ctx, TbCoreQueueFactory tbCoreQueueFactory, StreamObserver outputStream, BiConsumer sessionOpenListener, - BiConsumer sessionCloseListener, - ScheduledExecutorService sendDownlinkExecutorService, + BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); - } - - protected void initConsumer(Supplier>> edgeEventsConsumer, long schedulerPoolSize) { - this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); - this.highPriorityExecutorService = Executors.newScheduledThreadPool((int) schedulerPoolSize, ThingsBoardThreadFactory.forName("edge-event-high-priority-scheduler")); - this.consumer = QueueConsumerManager.>builder() - .name("TB Edge events") - .msgPackProcessor(this::processMsgs) - .pollInterval(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval()) - .consumerCreator(edgeEventsConsumer) - .consumerExecutor(consumerExecutor) - .threadPrefix("edge-events") - .build(); - scheduleCheckForHighPriorityEvent(); + edgeEventsConsumer = tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edge.getId()); } private void processMsgs(List> msgs, TbQueueConsumer> consumer) { log.trace("[{}][{}] starting processing edge events", tenantId, sessionId); - if (isConnected() && isSyncCompleted()) { - - if (!highPriorityQueue.isEmpty()) { - processHighPriorityEvents(); - } else { - List edgeEvents = new ArrayList<>(); - for (TbProtoQueueMsg msg : msgs) { - EdgeEvent edgeEvent = ProtoUtils.fromProto(msg.getValue().getEdgeEventMsg()); - edgeEvents.add(edgeEvent); - } - List downlinkMsgsPack = convertToDownlinkMsgsPack(edgeEvents); - try { - boolean isInterrupted = sendDownlinkMsgsPack(downlinkMsgsPack).get(); - if (isInterrupted) { - log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId); - } else { - consumer.commit(); - } - } catch (Exception e) { - log.error("[{}] Failed to process all downlink messages", sessionId, e); - } + if (isConnected() && isSyncCompleted() && !isHighPriorityProcessing) { + List edgeEvents = new ArrayList<>(); + for (TbProtoQueueMsg msg : msgs) { + EdgeEvent edgeEvent = ProtoUtils.fromProto(msg.getValue().getEdgeEventMsg()); + edgeEvents.add(edgeEvent); } - } else { - log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId); - } - } - - private void scheduleCheckForHighPriorityEvent() { - highPriorityExecutorService.scheduleAtFixedRate(() -> { + List downlinkMsgsPack = convertToDownlinkMsgsPack(edgeEvents); try { - if (isConnected() && isSyncCompleted()) { - if (!highPriorityQueue.isEmpty()) { - processHighPriorityEvents(); - } + boolean isInterrupted = sendDownlinkMsgsPack(downlinkMsgsPack).get(); + if (isInterrupted) { + log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId); + } else { + consumer.commit(); } } catch (Exception e) { - log.error("Error in processing high priority events", e); + log.error("[{}] Failed to process all downlink messages", sessionId, e); } - }, 0, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval() * 3, TimeUnit.MILLISECONDS); - } - - public void startConsumers() { - consumer.subscribe(); - consumer.launch(); + } else { + log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId); + } } @PreDestroy @@ -130,4 +97,37 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession migrateEdgeEvents(boolean isMigrationProcessed) throws Exception { + if (isMigrationProcessed) { + return Futures.immediateFuture(Boolean.FALSE); + } + return super.processEdgeEvents(); + } + + @Override + public ListenableFuture processEdgeEvents() { + if (!isConsumerInit) { + this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); + this.consumer = QueueConsumerManager.>builder() + .name("TB Edge events") + .msgPackProcessor(this::processMsgs) + .pollInterval(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval()) + .consumerCreator(() -> edgeEventsConsumer) + .consumerExecutor(consumerExecutor) + .threadPrefix("edge-events") + .build(); + isConsumerInit = true; + consumer.subscribe(); + consumer.launch(); + } + return Futures.immediateFuture(Boolean.FALSE); + } + + @Override + public void processHighPriorityEvents() { + isHighPriorityProcessing = true; + super.processHighPriorityEvents(); + } + } 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 4ef55e5961..d8f83c393a 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 @@ -15,6 +15,8 @@ */ package org.thingsboard.server.service.edge.rpc; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import io.grpc.stub.StreamObserver; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.edge.Edge; @@ -37,4 +39,9 @@ public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession migrateEdgeEvents(boolean isMigrationProcessed) { + 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 863be23e42..a613a2b5a1 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 @@ -17,7 +17,6 @@ package org.thingsboard.server.service.entitiy; import com.fasterxml.jackson.core.type.TypeReference; import jakarta.annotation.PostConstruct; -import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionalEventListener; @@ -35,6 +34,7 @@ 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; @@ -54,17 +54,34 @@ 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; -@Component -@RequiredArgsConstructor @Slf4j +@Component public class EntityStateSourcingListener { + private final TopicService topicService; private final TbClusterService tbClusterService; private final TenantService tenantService; + private final Optional kafkaSettings; + private final Optional kafkaTopicConfigs; + + public EntityStateSourcingListener(TopicService topicService, TbClusterService tbClusterService, TenantService tenantService, + Optional kafkaSettings, Optional kafkaTopicConfigs) { + this.topicService = topicService; + this.tbClusterService = tbClusterService; + this.tenantService = tenantService; + this.kafkaSettings = kafkaSettings; + this.kafkaTopicConfigs = kafkaTopicConfigs; + } + @PostConstruct public void init() { log.debug("EntityStateSourcingListener initiated"); @@ -137,7 +154,7 @@ public class EntityStateSourcingListener { log.debug("[{}][{}][{}] Handling entity deletion event: {}", tenantId, entityType, entityId, event); switch (entityType) { - case ASSET, ASSET_PROFILE, ENTITY_VIEW, CUSTOMER, EDGE, NOTIFICATION_RULE -> { + case ASSET, ASSET_PROFILE, ENTITY_VIEW, CUSTOMER, NOTIFICATION_RULE -> { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); } case NOTIFICATION_REQUEST -> { @@ -177,6 +194,10 @@ 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 -> {} } } @@ -247,6 +268,14 @@ 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) {