|
|
@ -84,6 +84,7 @@ import java.util.concurrent.ConcurrentMap; |
|
|
import java.util.concurrent.ScheduledExecutorService; |
|
|
import java.util.concurrent.ScheduledExecutorService; |
|
|
import java.util.concurrent.ScheduledFuture; |
|
|
import java.util.concurrent.ScheduledFuture; |
|
|
import java.util.concurrent.TimeUnit; |
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
|
|
import java.util.concurrent.atomic.AtomicBoolean; |
|
|
import java.util.concurrent.locks.Lock; |
|
|
import java.util.concurrent.locks.Lock; |
|
|
import java.util.concurrent.locks.ReentrantLock; |
|
|
import java.util.concurrent.locks.ReentrantLock; |
|
|
import java.util.function.Consumer; |
|
|
import java.util.function.Consumer; |
|
|
@ -722,6 +723,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void scheduleDisconnectNotification(Edge edge) { |
|
|
private void scheduleDisconnectNotification(Edge edge) { |
|
|
|
|
|
// Zero delay means "no debounce": notify immediately and skip the cluster-wide re-verify guard.
|
|
|
if (disconnectNotificationDelayMs <= 0) { |
|
|
if (disconnectNotificationDelayMs <= 0) { |
|
|
notifyEdgeConnectivity(edge, false); |
|
|
notifyEdgeConnectivity(edge, false); |
|
|
return; |
|
|
return; |
|
|
@ -729,7 +731,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i |
|
|
EdgeId edgeId = edge.getId(); |
|
|
EdgeId edgeId = edge.getId(); |
|
|
pendingDisconnectNotifications.compute(edgeId, (id, existing) -> { |
|
|
pendingDisconnectNotifications.compute(edgeId, (id, existing) -> { |
|
|
if (existing != null) { |
|
|
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
|
|
|
// 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()
|
|
|
// 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) { |
|
|
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(); |
|
|
EdgeId edgeId = edge.getId(); |
|
|
// Identity-keyed remove: don't clobber a newer entry if a second disconnect races with this task firing.
|
|
|
// Identity-keyed remove: don't clobber a newer entry if a second disconnect races with this task firing.
|
|
|
pendingDisconnectNotifications.remove(edgeId, pending); |
|
|
pendingDisconnectNotifications.remove(edgeId, pending); |
|
|
@ -760,23 +768,25 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i |
|
|
private void cancelPendingDisconnectNotification(EdgeId edgeId) { |
|
|
private void cancelPendingDisconnectNotification(EdgeId edgeId) { |
|
|
PendingDisconnect pending = pendingDisconnectNotifications.remove(edgeId); |
|
|
PendingDisconnect pending = pendingDisconnectNotifications.remove(edgeId); |
|
|
if (pending != null) { |
|
|
if (pending != null) { |
|
|
cancelIfPending(pending.future()); |
|
|
cancelIfPending(pending.getFuture()); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
static final class PendingDisconnect { |
|
|
static final class PendingDisconnect { |
|
|
|
|
|
|
|
|
private final Edge edge; |
|
|
private final Edge edge; |
|
|
|
|
|
private final AtomicBoolean notified = new AtomicBoolean(false); |
|
|
private ScheduledFuture<?> future; |
|
|
private ScheduledFuture<?> future; |
|
|
|
|
|
|
|
|
PendingDisconnect(Edge edge) { |
|
|
PendingDisconnect(Edge edge) { |
|
|
this.edge = edge; |
|
|
this.edge = edge; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
Edge edge() { |
|
|
Edge getEdge() { |
|
|
return edge; |
|
|
return edge; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
ScheduledFuture<?> future() { |
|
|
ScheduledFuture<?> getFuture() { |
|
|
return future; |
|
|
return future; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -784,6 +794,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i |
|
|
this.future = future; |
|
|
this.future = future; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
boolean tryClaim() { |
|
|
|
|
|
return notified.compareAndSet(false, true); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private static void cancelIfPending(ScheduledFuture<?> future) { |
|
|
private static void cancelIfPending(ScheduledFuture<?> future) { |
|
|
|