From 668ca3b632674bc42d1616d4de58d18c5e94be26 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 11 Jun 2026 10:55:51 +0300 Subject: [PATCH] Fixes after dev-rereview --- .../service/edge/rpc/EdgeGrpcService.java | 24 +++++++++++++++---- .../edge/EdgeConnectionNotificationTest.java | 9 +++++-- 2 files changed, 26 insertions(+), 7 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 411ccd5f05..297a45ec2f 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 @@ -84,6 +84,7 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; @@ -722,6 +723,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void scheduleDisconnectNotification(Edge edge) { + // Zero delay means "no debounce": notify immediately and skip the cluster-wide re-verify guard. if (disconnectNotificationDelayMs <= 0) { notifyEdgeConnectivity(edge, false); return; @@ -729,7 +731,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i EdgeId edgeId = edge.getId(); pendingDisconnectNotifications.compute(edgeId, (id, existing) -> { if (existing != null) { - cancelIfPending(existing.future()); + cancelIfPending(existing.getFuture()); } // Create the pending entry first so the scheduled task can reference it (for the identity-keyed // remove in fireDelayedDisconnectNotification), then back-fill its future. Doing this inside compute() @@ -744,7 +746,13 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void fireDelayedDisconnectNotification(PendingDisconnect pending) { - Edge edge = pending.edge(); + // Claim-once guard: the @PreDestroy destroy() flush and a concurrently-firing scheduled task can both call + // this for the same PendingDisconnect. tryClaim() ensures notifyEdgeConnectivity fires at most once without + // relying on downstream notification dedup. + if (!pending.tryClaim()) { + return; + } + Edge edge = pending.getEdge(); EdgeId edgeId = edge.getId(); // Identity-keyed remove: don't clobber a newer entry if a second disconnect races with this task firing. pendingDisconnectNotifications.remove(edgeId, pending); @@ -760,23 +768,25 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void cancelPendingDisconnectNotification(EdgeId edgeId) { PendingDisconnect pending = pendingDisconnectNotifications.remove(edgeId); if (pending != null) { - cancelIfPending(pending.future()); + cancelIfPending(pending.getFuture()); } } static final class PendingDisconnect { + private final Edge edge; + private final AtomicBoolean notified = new AtomicBoolean(false); private ScheduledFuture future; PendingDisconnect(Edge edge) { this.edge = edge; } - Edge edge() { + Edge getEdge() { return edge; } - ScheduledFuture future() { + ScheduledFuture getFuture() { return future; } @@ -784,6 +794,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i this.future = future; } + boolean tryClaim() { + return notified.compareAndSet(false, true); + } + } private static void cancelIfPending(ScheduledFuture future) { diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeConnectionNotificationTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeConnectionNotificationTest.java index 66023c7372..4482177750 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EdgeConnectionNotificationTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/EdgeConnectionNotificationTest.java @@ -30,6 +30,7 @@ import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.edge.imitator.EdgeImitator; import org.thingsboard.server.service.edge.rpc.EdgeGrpcService; +import java.util.Map; import java.util.concurrent.TimeUnit; import static org.awaitility.Awaitility.await; @@ -101,8 +102,12 @@ public class EdgeConnectionNotificationTest extends AbstractEdgeTest { // Edge drops... edgeImitator.disconnect(); - // Ensure the server processed the disconnect (and scheduled the delayed notification) before reconnecting. - TimeUnit.SECONDS.sleep(1); + // Wait until the server has processed the disconnect and scheduled the pending notification + // (the edge id appears in the pendingDisconnectNotifications map) before reconnecting. + await().atMost(AbstractWebTest.TIMEOUT, TimeUnit.SECONDS).until(() -> { + Map pending = (Map) ReflectionTestUtils.getField(edgeGrpcService, "pendingDisconnectNotifications"); + return pending != null && pending.containsKey(edge.getId()); + }); // ...and reconnects within the delay window, which must cancel the pending "disconnected" notification. EdgeImitator reconnected = createEdgeImitator();