From e4ad8877a31f0853ca1ab207bffda50a34b917f6 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 25 Nov 2025 18:39:48 +0200 Subject: [PATCH] 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;