From e4ad8877a31f0853ca1ab207bffda50a34b917f6 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 25 Nov 2025 18:39:48 +0200 Subject: [PATCH 1/2] Edge Zombie session fix - added session by session id map to handle properly connect and disconnect edge events --- .../service/edge/rpc/EdgeGrpcService.java | 39 ++++++++++++++++--- .../edge/rpc/KafkaEdgeGrpcSession.java | 20 +++++++++- 2 files changed, 53 insertions(+), 6 deletions(-) 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 ec3f839cb3..5159e17dbe 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 @@ -94,6 +94,7 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAS public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { private final ConcurrentMap sessions = new ConcurrentHashMap<>(); + private final ConcurrentMap sessionsById = new ConcurrentHashMap<>(); private final ConcurrentMap sessionNewEventsLocks = new ConcurrentHashMap<>(); private final Map sessionNewEvents = new HashMap<>(); private final ConcurrentMap> sessionEdgeEventChecks = new ConcurrentHashMap<>(); @@ -283,6 +284,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i destroySession(session); session.cleanUp(); sessions.remove(edgeId); + sessionsById.remove(session.getSessionId()); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); newEventLock.lock(); try { @@ -332,9 +334,15 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i TenantId tenantId = edge.getTenantId(); log.info("[{}][{}] edge [{}] connected successfully.", tenantId, edgeGrpcSession.getSessionId(), edgeId); if (sessions.containsKey(edgeId)) { - destroySession(sessions.get(edgeId)); + EdgeGrpcSession existing = sessions.get(edgeId); + if (existing != null) { + log.info("[{}][{}] Replacing existing session [{}] for edge [{}]", tenantId, edgeGrpcSession.getSessionId(), existing.getSessionId(), edgeId); + destroySession(existing); + sessionsById.remove(existing.getSessionId()); + } } sessions.put(edgeId, edgeGrpcSession); + sessionsById.put(edgeGrpcSession.getSessionId(), edgeGrpcSession); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); newEventLock.lock(); try { @@ -492,9 +500,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void onEdgeDisconnect(Edge edge, UUID sessionId) { EdgeId edgeId = edge.getId(); log.info("[{}][{}] edge disconnected!", edgeId, sessionId); - EdgeGrpcSession toRemove = sessions.get(edgeId); - if (toRemove.getSessionId().equals(sessionId)) { - toRemove = sessions.remove(edgeId); + EdgeGrpcSession current = sessions.get(edgeId); + if (current != null && current.getSessionId().equals(sessionId)) { + EdgeGrpcSession toRemove = sessions.remove(edgeId); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); newEventLock.lock(); try { @@ -503,6 +511,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i newEventLock.unlock(); } destroySession(toRemove); + sessionsById.remove(sessionId); TenantId tenantId = toRemove.getEdge().getTenantId(); save(tenantId, edgeId, ACTIVITY_STATE, false); long lastDisconnectTs = System.currentTimeMillis(); @@ -510,7 +519,18 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i pushRuleEngineMessage(toRemove.getEdge().getTenantId(), edge, lastDisconnectTs, TbMsgType.DISCONNECT_EVENT); cancelScheduleEdgeEventsCheck(edgeId); } else { - log.debug("[{}] edge session [{}] is not available anymore, nothing to remove. most probably this session is already outdated!", edgeId, sessionId); + log.info("[{}] edge session [{}] is not current anymore. Attempting to destroy it by sessionId.", edgeId, sessionId); + EdgeGrpcSession stale = sessionsById.remove(sessionId); + if (stale != null) { + try { + destroySession(stale); + log.info("[{}][{}] Successfully destroyed stale session for edge [{}]", stale.getTenantId(), sessionId, edgeId); + } catch (Exception e) { + log.warn("[{}][{}] Failed to destroy stale session for edge [{}]", stale.getTenantId(), sessionId, edgeId, e); + } + } else { + log.debug("[{}] No session found by sessionId [{}] to destroy", edgeId, sessionId); + } } edgeIdServiceIdCache.evict(edgeId); } @@ -522,6 +542,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i session.getTenantId(), session.getEdge().getId(), session.getEdge().getName(), session.getSessionId()); zombieSessions.add(session); } + } catch (Exception e) { + log.warn("[{}][{}] Exception during session destroy for edge [{}] with session id [{}]", + session.getTenantId(), session.getEdge().getId(), session.getEdge().getName(), session.getSessionId(), e); } } @@ -640,6 +663,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i !kafkaSession.getConsumer().getConsumer().isStopped()) { toRemove.add(kafkaSession.getEdge().getId()); } + if (session instanceof KafkaEdgeGrpcSession kafkaSession) { + log.debug("[{}] kafkaSession.isConnected() = {}, kafkaSession.getConsumer().getConsumer().isStopped() = {}", + kafkaSession.getEdge().getId(), + kafkaSession.isConnected(), + kafkaSession.getConsumer() != null ? kafkaSession.getConsumer().getConsumer() != null ? kafkaSession.getConsumer().getConsumer().isStopped() : null : null); + } } for (EdgeId edgeId : toRemove) { log.info("[{}] Destroying session for edge because edge is not connected", edgeId); 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 d165be33d4..67a5bd3623 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 @@ -101,8 +101,20 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { @Override public ListenableFuture processEdgeEvents() { + if (!isConnected() || isSyncInProgress() || isHighPriorityProcessing) { + log.warn("[{}][{}] Session is not ready (connected={}, syncInProgress={}, highPriority={}), skip starting edge event consumer", + tenantId, edge != null ? edge.getId() : null, isConnected(), isSyncInProgress(), isHighPriorityProcessing); + return Futures.immediateFuture(Boolean.FALSE); + } if (consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped())) { try { + if (this.consumerExecutor != null && !this.consumerExecutor.isShutdown()) { + try { + this.consumerExecutor.shutdown(); + } catch (Exception e) { + log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edge.getId(), e); + } + } this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); this.consumer = QueueConsumerManager.>builder() .name("TB Edge events [" + edge.getId() + "]") @@ -133,6 +145,7 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { public boolean destroy() { try { if (consumer != null) { + log.info("[{}][{}] Stopping edge event consumer...", tenantId, edge != null ? edge.getId() : null); consumer.stop(); } } catch (Exception e) { @@ -143,9 +156,14 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { try { if (consumerExecutor != null) { consumerExecutor.shutdown(); + try { + consumerExecutor.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS); + } catch (InterruptedException ie) { + log.warn("[{}][{}] Interrupted while awaiting consumer executor termination", tenantId, edge.getId()); + } } } catch (Exception e) { - log.warn("[{}][{}] Failed to shutdown consumer executor", tenantId, edge.getId(), e); + log.warn("[{}][{}] Failed to stop edge event consumer", tenantId, edge.getId(), e); return false; } return true; From 3a0d083610a07fd3ac07ba656f90abe4f7b6e07f Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 26 Nov 2025 10:18:05 +0200 Subject: [PATCH 2/2] Copilot codereview changes --- .../service/edge/rpc/EdgeGrpcService.java | 64 +++++++++++-------- .../edge/rpc/KafkaEdgeGrpcSession.java | 23 ++++--- 2 files changed, 53 insertions(+), 34 deletions(-) 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 5159e17dbe..c8b2c564ff 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 @@ -69,6 +69,7 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import java.io.IOException; import java.io.InputStream; import java.util.ArrayList; +import java.util.Collection; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -82,6 +83,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; +import java.util.function.Function; import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE; import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_CONNECT_TIME; @@ -654,31 +656,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void cleanupZombieSessions() { try { - List toRemove = new ArrayList<>(); - for (EdgeGrpcSession session : sessions.values()) { - if (session instanceof KafkaEdgeGrpcSession kafkaSession && - !kafkaSession.isConnected() && - kafkaSession.getConsumer() != null && - kafkaSession.getConsumer().getConsumer() != null && - !kafkaSession.getConsumer().getConsumer().isStopped()) { - toRemove.add(kafkaSession.getEdge().getId()); - } - if (session instanceof KafkaEdgeGrpcSession kafkaSession) { - log.debug("[{}] kafkaSession.isConnected() = {}, kafkaSession.getConsumer().getConsumer().isStopped() = {}", - kafkaSession.getEdge().getId(), - kafkaSession.isConnected(), - kafkaSession.getConsumer() != null ? kafkaSession.getConsumer().getConsumer() != null ? kafkaSession.getConsumer().getConsumer().isStopped() : null : null); - } - } - for (EdgeId edgeId : toRemove) { - log.info("[{}] Destroying session for edge because edge is not connected", edgeId); - EdgeGrpcSession removed = sessions.get(edgeId); - if (removed instanceof KafkaEdgeGrpcSession kafkaSession) { - if (kafkaSession.destroy()) { - sessions.remove(edgeId); - } - } - } + tryToDestroyZombieSessions(getZombieSessions(sessions.values()), s -> sessions.remove(s.getEdge().getId())); + tryToDestroyZombieSessions(getZombieSessions(sessionsById.values()), s -> sessionsById.remove(s.getSessionId())); + zombieSessions.removeIf(zombie -> { if (zombie.destroy()) { log.info("[{}][{}] Successfully cleaned up zombie session [{}] for edge [{}].", @@ -695,4 +675,38 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } + private List getZombieSessions(Collection sessions) { + List result = new ArrayList<>(); + for (EdgeGrpcSession session : sessions) { + if (isKafkaSessionAndZombie(session)) { + result.add(session); + } + } + return result; + } + + private void tryToDestroyZombieSessions(List sessionsToRemove, Function removeFunc) { + for (EdgeGrpcSession toRemove : sessionsToRemove) { + log.info("[{}] Destroying session for edge because edge is not connected", toRemove.getEdge().getId()); + if (toRemove.destroy()) { + removeFunc.apply(toRemove); + } + } + } + + private boolean isKafkaSessionAndZombie(EdgeGrpcSession session) { + if (session instanceof KafkaEdgeGrpcSession kafkaSession) { + log.debug("[{}] kafkaSession.isConnected() = {}, kafkaSession.getConsumer().getConsumer().isStopped() = {}", + kafkaSession.getEdge().getId(), + kafkaSession.isConnected(), + kafkaSession.getConsumer() != null ? kafkaSession.getConsumer().getConsumer() != null ? kafkaSession.getConsumer().getConsumer().isStopped() : null : null); + return !kafkaSession.isConnected() && + kafkaSession.getConsumer() != null && + kafkaSession.getConsumer().getConsumer() != null && + !kafkaSession.getConsumer().getConsumer().isStopped(); + } + return false; + + } + } 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 67a5bd3623..63669a8e3d 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 @@ -108,9 +108,10 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { } if (consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped())) { try { - if (this.consumerExecutor != null && !this.consumerExecutor.isShutdown()) { + if (consumerExecutor != null && !consumerExecutor.isShutdown()) { try { - this.consumerExecutor.shutdown(); + consumerExecutor.shutdown(); + awaitConsumerTermination(); } catch (Exception e) { log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edge.getId(), e); } @@ -154,21 +155,25 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { } consumer = null; try { - if (consumerExecutor != null) { + if (consumerExecutor != null && !consumerExecutor.isShutdown()) { consumerExecutor.shutdown(); - try { - consumerExecutor.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS); - } catch (InterruptedException ie) { - log.warn("[{}][{}] Interrupted while awaiting consumer executor termination", tenantId, edge.getId()); - } + awaitConsumerTermination(); } } catch (Exception e) { - log.warn("[{}][{}] Failed to stop edge event consumer", tenantId, edge.getId(), e); + log.warn("[{}][{}] Failed to shutdown edge event consumer executor", tenantId, edge.getId(), e); return false; } return true; } + private void awaitConsumerTermination() { + try { + consumerExecutor.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS); + } catch (InterruptedException ie) { + log.warn("[{}][{}] Interrupted while awaiting consumer executor termination", tenantId, edge.getId()); + } + } + @Override public void cleanUp() { String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edge.getId()).getTopic();