From f4d0aea5996dcf8c0592f109b228624847e44a14 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 8 Jul 2021 13:32:45 +0300 Subject: [PATCH] Added lock for new events updates --- .../service/edge/rpc/EdgeGrpcService.java | 56 +++++++++++++++---- 1 file changed, 45 insertions(+), 11 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 11948742bd..cdb649037b 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 @@ -50,6 +50,7 @@ import java.io.File; import java.io.IOException; import java.io.InputStream; import java.util.Collections; +import java.util.HashMap; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -59,6 +60,8 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; @Service @Slf4j @@ -67,7 +70,8 @@ import java.util.concurrent.TimeUnit; public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { private final ConcurrentMap sessions = new ConcurrentHashMap<>(); - private final ConcurrentMap sessionNewEvents = new ConcurrentHashMap<>(); + private final ConcurrentMap sessionNewEventsLocks = new ConcurrentHashMap<>(); + private final Map sessionNewEvents = new HashMap<>(); private final ConcurrentMap> sessionEdgeEventChecks = new ConcurrentHashMap<>(); private static final ObjectMapper mapper = new ObjectMapper(); @@ -176,7 +180,13 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); session.close(); sessions.remove(edgeId); - sessionNewEvents.remove(edgeId); + Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + lock.lock(); + try { + sessionNewEvents.remove(edgeId); + } finally { + lock.unlock(); + } cancelScheduleEdgeEventsCheck(edgeId); } } @@ -184,16 +194,28 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void onEdgeEvent(TenantId tenantId, EdgeId edgeId) { log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId.getId()); - if (Boolean.FALSE.equals(sessionNewEvents.get(edgeId))) { - log.trace("[{}] set session new events flag to true [{}]", tenantId, edgeId.getId()); - sessionNewEvents.put(edgeId, true); + Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + lock.lock(); + try { + if (Boolean.FALSE.equals(sessionNewEvents.get(edgeId))) { + log.trace("[{}] set session new events flag to true [{}]", tenantId, edgeId.getId()); + sessionNewEvents.put(edgeId, true); + } + } finally { + lock.unlock(); } } private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { log.info("[{}] edge [{}] connected successfully.", edgeGrpcSession.getSessionId(), edgeId); sessions.put(edgeId, edgeGrpcSession); - sessionNewEvents.put(edgeId, true); + Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + lock.lock(); + try { + sessionNewEvents.put(edgeId, true); + } finally { + lock.unlock(); + } save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true); save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis()); cancelScheduleEdgeEventsCheck(edgeId); @@ -217,10 +239,16 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (sessions.containsKey(edgeId)) { ScheduledFuture schedule = scheduler.schedule(() -> { try { - if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { - log.trace("[{}] Set session new events flag to false", edgeId.getId()); - sessionNewEvents.put(edgeId, false); - session.processEdgeEvents(); + Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + lock.lock(); + try { + if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { + log.trace("[{}] Set session new events flag to false", edgeId.getId()); + sessionNewEvents.put(edgeId, false); + session.processEdgeEvents(); + } + } finally { + lock.unlock(); } } catch (Exception e) { log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), e); @@ -249,7 +277,13 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void onEdgeDisconnect(EdgeId edgeId) { log.info("[{}] edge disconnected!", edgeId); sessions.remove(edgeId); - sessionNewEvents.remove(edgeId); + Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + lock.lock(); + try { + sessionNewEvents.remove(edgeId); + } finally { + lock.unlock(); + } save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis()); cancelScheduleEdgeEventsCheck(edgeId);