From 32fcfdac92f5fe29907d58195a15f9e1a4b01e47 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 6 Sep 2022 16:24:18 +0300 Subject: [PATCH] Moved edge rpc processing logic from tenant actor to edge service. Execute requests to edge service in a separate executor service --- .../server/actors/app/AppActor.java | 7 +- .../server/actors/tenant/TenantActor.java | 26 +----- .../service/edge/rpc/EdgeGrpcService.java | 89 ++++++++++++------- .../service/edge/rpc/EdgeRpcService.java | 10 +-- .../common/msg/edge/EdgeEventUpdateMsg.java | 4 +- .../common/msg/edge/EdgeSessionMsg.java | 24 +++++ .../common/msg/edge/FromEdgeSyncResponse.java | 4 +- .../common/msg/edge/ToEdgeSyncRequest.java | 4 +- 8 files changed, 97 insertions(+), 71 deletions(-) create mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeSessionMsg.java diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index 4583c78d11..eeb6d8e82e 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.aware.TenantAwareMsg; +import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.RuleEngineException; @@ -106,7 +107,7 @@ public class AppActor extends ContextAwareActor { case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG: case EDGE_SYNC_REQUEST_TO_EDGE_SESSION_MSG: case EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG: - onToTenantActorMsg((TenantAwareMsg) msg); + onToEdgeSessionMsg((EdgeSessionMsg) msg); break; case SESSION_TIMEOUT_MSG: ctx.broadcastToChildrenByType(msg, EntityType.TENANT); @@ -194,7 +195,7 @@ public class AppActor extends ContextAwareActor { () -> new TenantActor.ActorCreator(systemContext, tenantId)); } - private void onToTenantActorMsg(TenantAwareMsg msg) { + private void onToEdgeSessionMsg(EdgeSessionMsg msg) { TbActorRef target = null; if (ModelConstants.SYSTEM_TENANT.equals(msg.getTenantId())) { log.warn("Message has system tenant id: {}", msg); @@ -204,7 +205,7 @@ public class AppActor extends ContextAwareActor { if (target != null) { target.tellWithHighPriority(msg); } else { - log.debug("[{}] Invalid edge event update msg: {}", msg.getTenantId(), msg); + log.debug("[{}] Invalid edge session msg: {}", msg.getTenantId(), msg); } } diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index 5906b73db8..5718f581ad 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -47,9 +47,7 @@ import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; -import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg; -import org.thingsboard.server.common.msg.edge.FromEdgeSyncResponse; -import org.thingsboard.server.common.msg.edge.ToEdgeSyncRequest; +import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; @@ -171,7 +169,7 @@ public class TenantActor extends RuleChainManagerActor { case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG: case EDGE_SYNC_REQUEST_TO_EDGE_SESSION_MSG: case EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG: - onToEdgeSessionMsg(msg); + onToEdgeSessionMsg((EdgeSessionMsg) msg); break; default: return false; @@ -275,24 +273,8 @@ public class TenantActor extends RuleChainManagerActor { () -> new DeviceActorCreator(systemContext, tenantId, deviceId)); } - private void onToEdgeSessionMsg(TbActorMsg msg) { - switch (msg.getMsgType()) { - case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG: - EdgeEventUpdateMsg edgeEventUpdateMsg = (EdgeEventUpdateMsg) msg; - log.trace("[{}] onToEdgeSessionMsg [{}]", edgeEventUpdateMsg.getTenantId(), msg); - systemContext.getEdgeRpcService().onEdgeEvent(tenantId, edgeEventUpdateMsg.getEdgeId()); - break; - case EDGE_SYNC_REQUEST_TO_EDGE_SESSION_MSG: - ToEdgeSyncRequest toEdgeSyncRequest = (ToEdgeSyncRequest) msg; - log.trace("[{}] toEdgeSyncRequest [{}]", toEdgeSyncRequest.getTenantId(), msg); - systemContext.getEdgeRpcService().startSyncProcess(tenantId, toEdgeSyncRequest.getEdgeId(), toEdgeSyncRequest.getId()); - break; - case EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG: - FromEdgeSyncResponse fromEdgeSyncResponse = (FromEdgeSyncResponse) msg; - log.trace("[{}] fromEdgeSyncResponse [{}]", fromEdgeSyncResponse.getTenantId(), msg); - systemContext.getEdgeRpcService().processSyncResponse(fromEdgeSyncResponse); - break; - } + private void onToEdgeSessionMsg(EdgeSessionMsg msg) { + systemContext.getEdgeRpcService().onToEdgeSessionMsg(tenantId, msg); } private ApiUsageState getApiUsageState() { 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 18d53fea71..70abad7ba6 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 @@ -35,6 +35,8 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg; +import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; import org.thingsboard.server.common.msg.edge.FromEdgeSyncResponse; import org.thingsboard.server.common.msg.edge.ToEdgeSyncRequest; import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc; @@ -113,7 +115,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private ScheduledExecutorService sendDownlinkExecutorService; - private ScheduledExecutorService syncScheduler; + private ScheduledExecutorService executorService; @PostConstruct public void init() { @@ -140,9 +142,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i log.error("Failed to start Edge RPC server!", e); throw new RuntimeException("Failed to start Edge RPC server!"); } - this.edgeEventProcessingExecutorService = Executors.newScheduledThreadPool(schedulerPoolSize, ThingsBoardThreadFactory.forName("edge-scheduler")); + this.edgeEventProcessingExecutorService = Executors.newScheduledThreadPool(schedulerPoolSize, ThingsBoardThreadFactory.forName("edge-event-check-scheduler")); this.sendDownlinkExecutorService = Executors.newScheduledThreadPool(sendSchedulerPoolSize, ThingsBoardThreadFactory.forName("edge-send-scheduler")); - this.syncScheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("edge-sync-scheduler")); + this.executorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("edge-service")); log.info("Edge RPC service initialized!"); } @@ -165,6 +167,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (sendDownlinkExecutorService != null) { sendDownlinkExecutorService.shutdownNow(); } + if (executorService != null) { + executorService.shutdownNow(); + } } @Override @@ -172,37 +177,63 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService).getInputStream(); } + @Override + public void onToEdgeSessionMsg(TenantId tenantId, EdgeSessionMsg msg) { + executorService.execute(() -> { + switch (msg.getMsgType()) { + case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG: + EdgeEventUpdateMsg edgeEventUpdateMsg = (EdgeEventUpdateMsg) msg; + log.trace("[{}] onToEdgeSessionMsg [{}]", edgeEventUpdateMsg.getTenantId(), msg); + onEdgeEvent(tenantId, edgeEventUpdateMsg.getEdgeId()); + break; + case EDGE_SYNC_REQUEST_TO_EDGE_SESSION_MSG: + ToEdgeSyncRequest toEdgeSyncRequest = (ToEdgeSyncRequest) msg; + log.trace("[{}] toEdgeSyncRequest [{}]", toEdgeSyncRequest.getTenantId(), msg); + startSyncProcess(tenantId, toEdgeSyncRequest.getEdgeId(), toEdgeSyncRequest.getId()); + break; + case EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG: + FromEdgeSyncResponse fromEdgeSyncResponse = (FromEdgeSyncResponse) msg; + log.trace("[{}] fromEdgeSyncResponse [{}]", fromEdgeSyncResponse.getTenantId(), msg); + processSyncResponse(fromEdgeSyncResponse); + break; + } + }); + } + @Override public void updateEdge(TenantId tenantId, Edge edge) { - EdgeGrpcSession session = sessions.get(edge.getId()); - if (session != null && session.isConnected()) { - log.debug("[{}] Updating configuration for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); - session.onConfigurationUpdate(edge); - } else { - log.debug("[{}] Session doesn't exist for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); - } + executorService.execute(() -> { + EdgeGrpcSession session = sessions.get(edge.getId()); + if (session != null && session.isConnected()) { + log.debug("[{}] Updating configuration for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); + session.onConfigurationUpdate(edge); + } else { + log.debug("[{}] Session doesn't exist for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); + } + }); } @Override public void deleteEdge(TenantId tenantId, EdgeId edgeId) { - EdgeGrpcSession session = sessions.get(edgeId); - if (session != null && session.isConnected()) { - log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); - session.close(); - sessions.remove(edgeId); - final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); - newEventLock.lock(); - try { - sessionNewEvents.remove(edgeId); - } finally { - newEventLock.unlock(); + executorService.execute(() -> { + EdgeGrpcSession session = sessions.get(edgeId); + if (session != null && session.isConnected()) { + log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); + session.close(); + sessions.remove(edgeId); + final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); + newEventLock.lock(); + try { + sessionNewEvents.remove(edgeId); + } finally { + newEventLock.unlock(); + } + cancelScheduleEdgeEventsCheck(edgeId); } - cancelScheduleEdgeEventsCheck(edgeId); - } + }); } - @Override - public void onEdgeEvent(TenantId tenantId, EdgeId edgeId) { + private void onEdgeEvent(TenantId tenantId, EdgeId edgeId) { EdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId.getId()); @@ -235,8 +266,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i scheduleEdgeEventsCheck(edgeGrpcSession); } - @Override - public void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId) { + private void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId) { EdgeGrpcSession session = sessions.get(edgeId); if (session != null) { boolean success = false; @@ -259,7 +289,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void scheduleSyncRequestTimeout(ToEdgeSyncRequest request, UUID requestId) { log.trace("[{}] scheduling sync edge request", requestId); - syncScheduler.schedule(() -> { + executorService.schedule(() -> { log.trace("[{}] checking if sync edge request is not processed...", requestId); Consumer consumer = localSyncEdgeRequests.remove(requestId); if (consumer != null) { @@ -269,8 +299,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i }, 10, TimeUnit.SECONDS); } - @Override - public void processSyncResponse(FromEdgeSyncResponse response) { + private void processSyncResponse(FromEdgeSyncResponse response) { log.trace("[{}] Received response from sync service: [{}]", response.getId(), response); UUID requestId = response.getId(); Consumer consumer = localSyncEdgeRequests.remove(requestId); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeRpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeRpcService.java index 9564e195de..6c5515337e 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeRpcService.java @@ -18,23 +18,19 @@ package org.thingsboard.server.service.edge.rpc; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; import org.thingsboard.server.common.msg.edge.FromEdgeSyncResponse; import org.thingsboard.server.common.msg.edge.ToEdgeSyncRequest; -import java.util.UUID; import java.util.function.Consumer; public interface EdgeRpcService { + void onToEdgeSessionMsg(TenantId tenantId, EdgeSessionMsg msg); + void updateEdge(TenantId tenantId, Edge edge); void deleteEdge(TenantId tenantId, EdgeId edgeId); - void onEdgeEvent(TenantId tenantId, EdgeId edgeId); - - void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId); - void processSyncRequest(ToEdgeSyncRequest request, Consumer responseConsumer); - - void processSyncResponse(FromEdgeSyncResponse response); } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java index 3469cde65f..0285eef0b1 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java @@ -20,11 +20,9 @@ import lombok.ToString; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; -import org.thingsboard.server.common.msg.aware.TenantAwareMsg; -import org.thingsboard.server.common.msg.cluster.ToAllNodesMsg; @ToString -public class EdgeEventUpdateMsg implements TenantAwareMsg, ToAllNodesMsg { +public class EdgeEventUpdateMsg implements EdgeSessionMsg { @Getter private final TenantId tenantId; @Getter diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeSessionMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeSessionMsg.java new file mode 100644 index 0000000000..c4719f9b29 --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeSessionMsg.java @@ -0,0 +1,24 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.msg.edge; + +import org.thingsboard.server.common.msg.aware.TenantAwareMsg; +import org.thingsboard.server.common.msg.cluster.ToAllNodesMsg; + +import java.io.Serializable; + +public interface EdgeSessionMsg extends TenantAwareMsg, ToAllNodesMsg { +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/FromEdgeSyncResponse.java b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/FromEdgeSyncResponse.java index 4e6d456d22..e93e5cd3b5 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/FromEdgeSyncResponse.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/FromEdgeSyncResponse.java @@ -20,14 +20,12 @@ import lombok.Getter; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; -import org.thingsboard.server.common.msg.aware.TenantAwareMsg; -import org.thingsboard.server.common.msg.cluster.ToAllNodesMsg; import java.util.UUID; @AllArgsConstructor @Getter -public class FromEdgeSyncResponse implements TenantAwareMsg, ToAllNodesMsg { +public class FromEdgeSyncResponse implements EdgeSessionMsg { private final UUID id; private final TenantId tenantId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/ToEdgeSyncRequest.java b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/ToEdgeSyncRequest.java index 6e1b0df3ad..32e1068e73 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/edge/ToEdgeSyncRequest.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/edge/ToEdgeSyncRequest.java @@ -20,14 +20,12 @@ import lombok.Getter; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; -import org.thingsboard.server.common.msg.aware.TenantAwareMsg; -import org.thingsboard.server.common.msg.cluster.ToAllNodesMsg; import java.util.UUID; @AllArgsConstructor @Getter -public class ToEdgeSyncRequest implements TenantAwareMsg, ToAllNodesMsg { +public class ToEdgeSyncRequest implements EdgeSessionMsg { private final UUID id; private final TenantId tenantId; private final EdgeId edgeId;