diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index d3da4f7f15..065c3f2cec 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -99,15 +99,13 @@ import java.util.concurrent.ScheduledExecutorService; @TbCoreComponent public class EdgeContextComponent { - @Value("${edges.scheduler_pool_size}") - private int schedulerPoolSize; - private final Map processorMap = new EnumMap<>(EdgeEventType.class); private ScheduledExecutorService edgeEventProcessingExecutorService; @Autowired - public EdgeContextComponent(List processors) { + public EdgeContextComponent(List processors, + @Value("${edges.scheduler_pool_size}") int schedulerPoolSize) { processors.forEach(processor -> { EdgeEventType eventType = processor.getEdgeEventType(); if (eventType != null) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java deleted file mode 100644 index 7d0f14c24e..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java +++ /dev/null @@ -1,43 +0,0 @@ -/** - * Copyright © 2016-2026 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.service.edge.rpc; - -import com.google.common.util.concurrent.FutureCallback; -import jakarta.annotation.Nullable; -import lombok.AllArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.id.EdgeId; -import org.thingsboard.server.common.data.id.TenantId; - -@Slf4j -@AllArgsConstructor -public class AttributeSaveCallback implements FutureCallback { - - private final TenantId tenantId; - private final EdgeId edgeId; - private final String key; - private final Object value; - - @Override - public void onSuccess(@Nullable Void result) { - log.trace("[{}][{}] Successfully updated attribute [{}] with value [{}]", tenantId, edgeId, key, value); - } - - @Override - public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to update attribute [{}] with value [{}]", tenantId, edgeId, key, value, t); - } -} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/DownlinkMessageMapper.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/DownlinkMessageMapper.java index 5c6d675f69..a83e3036f4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/DownlinkMessageMapper.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/DownlinkMessageMapper.java @@ -17,12 +17,15 @@ package org.thingsboard.server.service.edge.rpc; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.EdgeVersion; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeContextComponent; +import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; import org.thingsboard.server.service.edge.rpc.utils.EdgeVersionUtils; import java.util.ArrayList; @@ -31,13 +34,16 @@ import java.util.List; @Component @Slf4j @RequiredArgsConstructor +@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true") +@TbCoreComponent public class DownlinkMessageMapper { private final EdgeContextComponent ctx; public List convertToDownlinkMsgsPack(EdgeSessionState state, List edgeEvents) { List result = new ArrayList<>(); - for (EdgeEvent edgeEvent : edgeEvents) { + List filtered = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(edgeEvents); + for (EdgeEvent edgeEvent : filtered) { log.trace("[{}][{}] converting edge event to downlink msg [{}]", state.getTenantId(), state.getEdgeId(), edgeEvent); DownlinkMsg downlinkMsg = null; try { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeUplinkMessageDispatcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeUplinkMessageDispatcher.java index 852a6f2ad7..e5dd0738ac 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeUplinkMessageDispatcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeUplinkMessageDispatcher.java @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.notification.rule.trigger.EdgeCommunicationFailureTrigger; @@ -51,6 +52,7 @@ import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeContextComponent; import java.util.ArrayList; @@ -59,6 +61,8 @@ import java.util.List; @Service @Slf4j @RequiredArgsConstructor +@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true") +@TbCoreComponent public class EdgeUplinkMessageDispatcher { private final EdgeContextComponent ctx; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/GrpcServer.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/GrpcServer.java index 4ca8a4f838..c5c2e40f5b 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/GrpcServer.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/GrpcServer.java @@ -24,11 +24,13 @@ import jakarta.annotation.PreDestroy; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.transport.config.ssl.PemSslCredentials; import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc.EdgeRpcServiceImplBase; import org.thingsboard.server.queue.util.AfterStartUp; +import org.thingsboard.server.queue.util.TbCoreComponent; import java.io.IOException; import java.util.concurrent.TimeUnit; @@ -36,6 +38,8 @@ import java.util.concurrent.TimeUnit; @Component @Slf4j @RequiredArgsConstructor +@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true") +@TbCoreComponent public class GrpcServer { @Value("${edges.rpc.port}") diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java index f9ecac9923..a41ff78110 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java @@ -23,6 +23,7 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; @@ -54,6 +55,7 @@ import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc; import org.thingsboard.server.gen.edge.v1.RequestMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.edge.rpc.EdgeSessionState; @@ -77,6 +79,8 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAS @Service @Slf4j @RequiredArgsConstructor +@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true") +@TbCoreComponent public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { @Value("${edges.send_scheduler_pool_size}") diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/DefaultZombieSessionCleanupService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/DefaultZombieSessionCleanupService.java index 6aac5646d8..76d554eab9 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/DefaultZombieSessionCleanupService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/DefaultZombieSessionCleanupService.java @@ -18,8 +18,10 @@ package org.thingsboard.server.service.edge.rpc.session; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardExecutors; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.KafkaBasedEdgeGrpcSessionManager; @@ -35,6 +37,8 @@ import java.util.function.Function; @Service @Slf4j +@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true") +@TbCoreComponent public class DefaultZombieSessionCleanupService implements ZombieSessionCleanupService { @Autowired diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java index f6705b13e2..dc2296f1d5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java @@ -58,7 +58,6 @@ import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; -import org.thingsboard.server.service.edge.rpc.AttributeSaveCallback; import org.thingsboard.server.service.edge.rpc.DownlinkMessageMapper; import org.thingsboard.server.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.EdgeSyncCursor; @@ -604,7 +603,7 @@ public class EdgeGrpcSession implements EdgeSession { .entityId(getEdgeId()) .scope(AttributeScope.SERVER_SCOPE) .entry(new BooleanDataEntry(DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY, value)) - .callback(new AttributeSaveCallback(getTenantId(), getEdgeId(), DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY, value)) + .callback(new EdgeAttributeSaveCallback(getTenantId(), getEdgeId(), DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY, value)) .build()); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java index 332da77ad6..7cc27214a7 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java @@ -17,8 +17,10 @@ package org.thingsboard.server.service.edge.rpc.session; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; @@ -30,6 +32,8 @@ import java.util.function.Consumer; @Data @Slf4j @Component +@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true") +@TbCoreComponent public class EdgeSessionsHolder { private final ConcurrentMap sessions = new ConcurrentHashMap<>(); @@ -66,6 +70,7 @@ public class EdgeSessionsHolder { public void remove(EdgeGrpcSessionManager session) { if (session == null) { log.warn("Can't remove session from holder because it's null"); + return; } EdgeSessionState sessionState = session.getState(); removeByEdgeId(sessionState.getEdge().getId()); // todo: react to warnings diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java index 91c33cfc6d..186510d5bb 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java @@ -167,9 +167,10 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan try { isHighPriorityProcessing = true; session.processHighPriorityEvents(); - isHighPriorityProcessing = false; } catch (Exception e) { log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, edgeId, e); + } finally { + isHighPriorityProcessing = false; } }, NO_INITIAL_DELAY_VALUE, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS); highPriorityProcessingFutureRef.set(highPriorityProcessingTask); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/PostgresBasedEdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/PostgresBasedEdgeGrpcSessionManager.java index 5b86e309e6..242f513e10 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/PostgresBasedEdgeGrpcSessionManager.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/PostgresBasedEdgeGrpcSessionManager.java @@ -168,7 +168,7 @@ public class PostgresBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSession private void markHasEvents(boolean newEventsPresent) { newEventsLock.lock(); try { - if (!hasNewEvents) { + if (hasNewEvents != newEventsPresent) { log.trace("[{}] set session new events flag to {} [{}]", getState().getTenantId(), newEventsPresent, getState().getEdgeId()); hasNewEvents = newEventsPresent; }