Browse Source

Code review changes

pull/14618/head
Volodymyr Babak 7 months ago
parent
commit
06bf79234a
  1. 6
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 43
      application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java
  3. 8
      application/src/main/java/org/thingsboard/server/service/edge/rpc/DownlinkMessageMapper.java
  4. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeUplinkMessageDispatcher.java
  5. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/GrpcServer.java
  6. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java
  7. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/DefaultZombieSessionCleanupService.java
  8. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java
  9. 5
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java
  10. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java
  11. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/PostgresBasedEdgeGrpcSessionManager.java

6
application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java

@ -99,15 +99,13 @@ import java.util.concurrent.ScheduledExecutorService;
@TbCoreComponent @TbCoreComponent
public class EdgeContextComponent { public class EdgeContextComponent {
@Value("${edges.scheduler_pool_size}")
private int schedulerPoolSize;
private final Map<EdgeEventType, EdgeProcessor> processorMap = new EnumMap<>(EdgeEventType.class); private final Map<EdgeEventType, EdgeProcessor> processorMap = new EnumMap<>(EdgeEventType.class);
private ScheduledExecutorService edgeEventProcessingExecutorService; private ScheduledExecutorService edgeEventProcessingExecutorService;
@Autowired @Autowired
public EdgeContextComponent(List<EdgeProcessor> processors) { public EdgeContextComponent(List<EdgeProcessor> processors,
@Value("${edges.scheduler_pool_size}") int schedulerPoolSize) {
processors.forEach(processor -> { processors.forEach(processor -> {
EdgeEventType eventType = processor.getEdgeEventType(); EdgeEventType eventType = processor.getEdgeEventType();
if (eventType != null) { if (eventType != null) {

43
application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java

@ -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<Void> {
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);
}
}

8
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.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg;
import org.thingsboard.server.gen.edge.v1.EdgeVersion; 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.EdgeContextComponent;
import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils;
import org.thingsboard.server.service.edge.rpc.utils.EdgeVersionUtils; import org.thingsboard.server.service.edge.rpc.utils.EdgeVersionUtils;
import java.util.ArrayList; import java.util.ArrayList;
@ -31,13 +34,16 @@ import java.util.List;
@Component @Component
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@TbCoreComponent
public class DownlinkMessageMapper { public class DownlinkMessageMapper {
private final EdgeContextComponent ctx; private final EdgeContextComponent ctx;
public List<DownlinkMsg> convertToDownlinkMsgsPack(EdgeSessionState state, List<EdgeEvent> edgeEvents) { public List<DownlinkMsg> convertToDownlinkMsgsPack(EdgeSessionState state, List<EdgeEvent> edgeEvents) {
List<DownlinkMsg> result = new ArrayList<>(); List<DownlinkMsg> result = new ArrayList<>();
for (EdgeEvent edgeEvent : edgeEvents) { List<EdgeEvent> filtered = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(edgeEvents);
for (EdgeEvent edgeEvent : filtered) {
log.trace("[{}][{}] converting edge event to downlink msg [{}]", state.getTenantId(), state.getEdgeId(), edgeEvent); log.trace("[{}][{}] converting edge event to downlink msg [{}]", state.getTenantId(), state.getEdgeId(), edgeEvent);
DownlinkMsg downlinkMsg = null; DownlinkMsg downlinkMsg = null;
try { try {

4
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 com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.notification.rule.trigger.EdgeCommunicationFailureTrigger; 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.UserCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; import org.thingsboard.server.gen.edge.v1.UserUpdateMsg;
import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg; import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.EdgeContextComponent;
import java.util.ArrayList; import java.util.ArrayList;
@ -59,6 +61,8 @@ import java.util.List;
@Service @Service
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@TbCoreComponent
public class EdgeUplinkMessageDispatcher { public class EdgeUplinkMessageDispatcher {
private final EdgeContextComponent ctx; private final EdgeContextComponent ctx;

4
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.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.transport.config.ssl.PemSslCredentials; import org.thingsboard.server.common.transport.config.ssl.PemSslCredentials;
import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc.EdgeRpcServiceImplBase; import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc.EdgeRpcServiceImplBase;
import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.queue.util.TbCoreComponent;
import java.io.IOException; import java.io.IOException;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@ -36,6 +38,8 @@ import java.util.concurrent.TimeUnit;
@Component @Component
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@TbCoreComponent
public class GrpcServer { public class GrpcServer {
@Value("${edges.rpc.port}") @Value("${edges.rpc.port}")

4
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 lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Lazy; import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service; 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.RequestMsg;
import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; 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.EdgeContextComponent;
import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.edge.rpc.EdgeRpcService;
import org.thingsboard.server.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.EdgeSessionState;
@ -77,6 +79,8 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAS
@Service @Service
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@TbCoreComponent
public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService {
@Value("${edges.send_scheduler_pool_size}") @Value("${edges.send_scheduler_pool_size}")

4
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 jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardExecutors; 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.EdgeSessionState;
import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager;
import org.thingsboard.server.service.edge.rpc.session.manager.KafkaBasedEdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.KafkaBasedEdgeGrpcSessionManager;
@ -35,6 +37,8 @@ import java.util.function.Function;
@Service @Service
@Slf4j @Slf4j
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@TbCoreComponent
public class DefaultZombieSessionCleanupService implements ZombieSessionCleanupService { public class DefaultZombieSessionCleanupService implements ZombieSessionCleanupService {
@Autowired @Autowired

3
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.gen.edge.v1.UplinkResponseMsg;
import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.EdgeContextComponent;
import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; 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.DownlinkMessageMapper;
import org.thingsboard.server.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.EdgeSessionState;
import org.thingsboard.server.service.edge.rpc.EdgeSyncCursor; import org.thingsboard.server.service.edge.rpc.EdgeSyncCursor;
@ -604,7 +603,7 @@ public class EdgeGrpcSession implements EdgeSession {
.entityId(getEdgeId()) .entityId(getEdgeId())
.scope(AttributeScope.SERVER_SCOPE) .scope(AttributeScope.SERVER_SCOPE)
.entry(new BooleanDataEntry(DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY, value)) .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()); .build());
} }

5
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.Data;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.EdgeId; 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.EdgeSessionState;
import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager;
@ -30,6 +32,8 @@ import java.util.function.Consumer;
@Data @Data
@Slf4j @Slf4j
@Component @Component
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@TbCoreComponent
public class EdgeSessionsHolder { public class EdgeSessionsHolder {
private final ConcurrentMap<EdgeId, EdgeGrpcSessionManager> sessions = new ConcurrentHashMap<>(); private final ConcurrentMap<EdgeId, EdgeGrpcSessionManager> sessions = new ConcurrentHashMap<>();
@ -66,6 +70,7 @@ public class EdgeSessionsHolder {
public void remove(EdgeGrpcSessionManager session) { public void remove(EdgeGrpcSessionManager session) {
if (session == null) { if (session == null) {
log.warn("Can't remove session from holder because it's null"); log.warn("Can't remove session from holder because it's null");
return;
} }
EdgeSessionState sessionState = session.getState(); EdgeSessionState sessionState = session.getState();
removeByEdgeId(sessionState.getEdge().getId()); // todo: react to warnings removeByEdgeId(sessionState.getEdge().getId()); // todo: react to warnings

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java

@ -167,9 +167,10 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
try { try {
isHighPriorityProcessing = true; isHighPriorityProcessing = true;
session.processHighPriorityEvents(); session.processHighPriorityEvents();
isHighPriorityProcessing = false;
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, edgeId, 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); }, NO_INITIAL_DELAY_VALUE, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS);
highPriorityProcessingFutureRef.set(highPriorityProcessingTask); highPriorityProcessingFutureRef.set(highPriorityProcessingTask);

2
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) { private void markHasEvents(boolean newEventsPresent) {
newEventsLock.lock(); newEventsLock.lock();
try { try {
if (!hasNewEvents) { if (hasNewEvents != newEventsPresent) {
log.trace("[{}] set session new events flag to {} [{}]", getState().getTenantId(), newEventsPresent, getState().getEdgeId()); log.trace("[{}] set session new events flag to {} [{}]", getState().getTenantId(), newEventsPresent, getState().getEdgeId());
hasNewEvents = newEventsPresent; hasNewEvents = newEventsPresent;
} }

Loading…
Cancel
Save