From 6fe25a2ee29ee68e78999050b1ba3265a1466aa7 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Sat, 28 Mar 2020 01:42:17 +0200 Subject: [PATCH] Refactoring to support TB-Rule-Engine service --- .../server/actors/ActorSystemContext.java | 26 +-- .../device/DeviceActorMessageProcessor.java | 23 +-- .../actors/ruleChain/RuleChainActor.java | 2 +- .../actors/ruleChain/RuleNodeActor.java | 2 + .../ThingsboardSecurityConfiguration.java | 1 + .../server/config/WebSocketConfiguration.java | 5 +- .../server/controller/AdminController.java | 14 +- .../server/controller/AlarmController.java | 2 + .../server/controller/AssetController.java | 2 + .../server/controller/AuditLogController.java | 2 + .../server/controller/AuthController.java | 2 + .../server/controller/BaseController.java | 3 + .../ComponentDescriptorController.java | 2 + .../server/controller/CustomerController.java | 2 + .../controller/DashboardController.java | 3 +- .../server/controller/DeviceController.java | 2 + .../controller/EntityRelationController.java | 4 +- .../controller/EntityViewController.java | 2 + .../server/controller/EventController.java | 2 + .../server/controller/RpcController.java | 3 +- .../controller/RuleChainController.java | 2 + .../controller/TelemetryController.java | 2 + .../server/controller/TenantController.java | 2 + .../server/controller/UserController.java | 2 + .../controller/WidgetTypeController.java | 2 + .../controller/WidgetsBundleController.java | 2 + .../controller/plugin/TbWebSocketHandler.java | 3 +- .../queue/DefaultTbCoreConsumerService.java | 155 ++++++------------ .../DefaultTbRuleEngineConsumerService.java | 154 +++++------------ .../processing/AbstractConsumerService.java | 150 +++++++++++++++++ .../rpc/DefaultTbCoreDeviceRpcService.java | 5 +- .../rpc/DefaultTbRuleEngineRpcService.java | 7 +- .../auth/rest/RestAuthenticationProvider.java | 2 + ...RestAwareAuthenticationFailureHandler.java | 1 + ...RestAwareAuthenticationSuccessHandler.java | 1 + .../device/DefaultDeviceAuthService.java | 13 +- .../security/model/token/JwtTokenFactory.java | 1 + .../permission/AccessControlService.java | 2 - .../DefaultAccessControlService.java | 1 + .../system/DefaultSystemSecurityService.java | 1 + .../DefaultDeviceSessionCacheService.java | 2 + .../state/DefaultDeviceStateService.java | 18 +- .../service/state/DeviceStateService.java | 2 +- .../DefaultSubscriptionManagerService.java | 2 + .../DefaultTbLocalSubscriptionService.java | 2 + .../DefaultTelemetrySubscriptionService.java | 47 ++++-- .../DefaultTelemetryWebSocketService.java | 2 + .../BaseRuleChainTransactionService.java | 10 +- .../DefaultTbCoreToTransportService.java | 3 +- .../transport/DefaultTransportApiService.java | 2 + ...ce.java => TbCoreTransportApiService.java} | 9 +- .../service/update/DefaultUpdateService.java | 2 + .../src/main/resources/thingsboard.yml | 11 +- .../server/common/msg/queue/ServiceKey.java | 2 + .../common/msg/queue/TopicPartitionInfo.java | 2 + .../provider/KafkaMonolithQueueProvider.java | 6 +- .../provider/KafkaTbCoreQueueProvider.java | 6 +- .../KafkaTbRuleEngineQueueProvider.java | 6 +- .../queue/provider/TbCoreQueueProvider.java | 8 +- .../server/queue/util/TbCoreComponent.java | 15 +- .../server/queue/util/TbKafkaQueue.java | 19 +-- .../queue/util/TbMonolithComponent.java | 19 +-- .../queue/util/TbMonolithOrCoreComponent.java | 22 +++ .../util/TbMonolithOrRuleEngineComponent.java | 22 +++ .../queue/util/TbRuleEngineComponent.java | 22 +++ .../transport/mqtt/MqttTransportHandler.java | 16 +- .../service/DefaultTransportService.java | 9 +- .../session/DeviceAwareSessionContext.java | 8 +- 68 files changed, 526 insertions(+), 380 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java rename application/src/main/java/org/thingsboard/server/service/transport/{RemoteTransportApiService.java => TbCoreTransportApiService.java} (93%) rename application/src/main/java/org/thingsboard/server/service/transport/ToTransportMsgEncoder.java => common/queue/src/main/java/org/thingsboard/server/queue/util/TbCoreComponent.java (60%) rename application/src/main/java/org/thingsboard/server/service/transport/ToRuleEngineMsgDecoder.java => common/queue/src/main/java/org/thingsboard/server/queue/util/TbKafkaQueue.java (53%) rename application/src/main/java/org/thingsboard/server/service/executors/ClusterRpcCallbackExecutorService.java => common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithComponent.java (53%) create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrCoreComponent.java create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrRuleEngineComponent.java create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/util/TbRuleEngineComponent.java diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 03b65e7ccc..bd7c25ddc9 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -70,7 +70,6 @@ import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.encoding.DataDecodingEncodingService; -import org.thingsboard.server.service.executors.ClusterRpcCallbackExecutorService; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.ExternalCallExecutorService; import org.thingsboard.server.service.executors.SharedEventLoopGroupService; @@ -125,10 +124,6 @@ public class ActorSystemContext { @Getter private DataDecodingEncodingService encodingService; - @Autowired - @Getter - private DeviceAuthService deviceAuthService; - @Autowired @Getter private DeviceService deviceService; @@ -212,10 +207,6 @@ public class ActorSystemContext { @Getter private MailExecutorService mailExecutor; - @Autowired - @Getter - private ClusterRpcCallbackExecutorService clusterRpcCallbackExecutor; - @Autowired @Getter private DbCallbackExecutorService dbCallbackExecutor; @@ -232,37 +223,34 @@ public class ActorSystemContext { @Getter private MailService mailService; - @Autowired + //TODO: separate context for TbCore and TbRuleEngine + @Autowired(required = false) @Getter private DeviceStateService deviceStateService; - @Autowired + @Autowired(required = false) @Getter private DeviceSessionCacheService deviceSessionCacheService; - @Lazy - @Autowired + @Autowired(required = false) @Getter private TbCoreToTransportService tbCoreToTransportService; /** * The following Service will be null if we operate in tb-core mode */ - @Lazy - @Getter @Autowired(required = false) + @Getter private TbRuleEngineDeviceRpcService tbRuleEngineDeviceRpcService; /** * The following Service will be null if we operate in tb-rule-engine mode */ - @Lazy - @Getter @Autowired(required = false) + @Getter private TbCoreDeviceRpcService tbCoreDeviceRpcService; - @Lazy - @Autowired + @Autowired(required = false) @Getter private RuleChainTransactionService ruleChainTransactionService; diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index aff6200199..307797f864 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -212,8 +212,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } void process(ActorContext context, TransportToDeviceActorMsgWrapper wrapper) { - //TODO 2.5 - boolean reportDeviceActivity = true; TransportToDeviceActorMsg msg = wrapper.getMsg(); TbMsgCallback callback = wrapper.getCallback(); if (msg.hasSessionEvent()) { @@ -225,37 +223,18 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (msg.hasSubscribeToRPC()) { processSubscriptionCommands(context, msg.getSessionInfo(), msg.getSubscribeToRPC()); } -// if (msg.hasPostAttributes()) { -// handlePostAttributesRequest(context, msg.getSessionInfo(), msg.getPostAttributes()); -// reportDeviceActivity = true; -// } -// if (msg.hasPostTelemetry()) { -// handlePostTelemetryRequest(context, msg.getSessionInfo(), msg.getPostTelemetry()); -// reportDeviceActivity = true; -// } if (msg.hasGetAttributes()) { handleGetAttributesRequest(context, msg.getSessionInfo(), msg.getGetAttributes()); } if (msg.hasToDeviceRPCCallResponse()) { processRpcResponses(context, msg.getSessionInfo(), msg.getToDeviceRPCCallResponse()); } -// if (msg.hasToServerRPCCallRequest()) { -// handleClientSideRPCRequest(context, msg.getSessionInfo(), msg.getToServerRPCCallRequest()); -// reportDeviceActivity = true; -// } if (msg.hasSubscriptionInfo()) { handleSessionActivity(context, msg.getSessionInfo(), msg.getSubscriptionInfo()); } - if (reportDeviceActivity) { - reportLogicalDeviceActivity(); - } callback.onSuccess(); } - private void reportLogicalDeviceActivity() { - systemContext.getDeviceStateService().onDeviceActivity(deviceId); - } - private void reportSessionOpen() { systemContext.getDeviceStateService().onDeviceConnect(deviceId); } @@ -469,6 +448,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (sessions.size() == 1) { reportSessionOpen(); } + systemContext.getDeviceStateService().onDeviceActivity(deviceId, System.currentTimeMillis()); dumpSessions(); } else if (msg.getEvent() == SessionEvent.CLOSED) { log.debug("[{}] Canceling subscriptions for closed session [{}]", deviceId, sessionId); @@ -496,6 +476,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (subscriptionInfo.getRpcSubscription()) { rpcSubscriptions.putIfAbsent(sessionId, sessionMD.getSessionInfo()); } + systemContext.getDeviceStateService().onDeviceActivity(deviceId, subscriptionInfo.getLastActivityTime()); dumpSessions(); } diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java index 1b9923139a..015e7a008b 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java @@ -53,7 +53,7 @@ public class RuleChainActor extends ComponentActor implements TbCoreConsumerService { @Value("${queue.core.poll_interval}") private long pollDuration; @@ -72,56 +73,33 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { @Value("${queue.core.stats.enabled:false}") private boolean statsEnabled; - private final ActorSystemContext actorContext; private final DeviceStateService stateService; private final TbLocalSubscriptionService localSubscriptionService; private final SubscriptionManagerService subscriptionManagerService; - private final DataDecodingEncodingService encodingService; - private final TbQueueConsumer> mainConsumer; - private final TbQueueConsumer> nfConsumer; private final TbCoreDeviceRpcService tbCoreDeviceRpcService; private final TbCoreConsumerStats stats = new TbCoreConsumerStats(); - private volatile ExecutorService mainConsumerExecutor; - private volatile ExecutorService notificationsConsumerExecutor; private volatile boolean stopped = false; public DefaultTbCoreConsumerService(TbCoreQueueProvider tbCoreQueueProvider, ActorSystemContext actorContext, DeviceStateService stateService, TbLocalSubscriptionService localSubscriptionService, SubscriptionManagerService subscriptionManagerService, DataDecodingEncodingService encodingService, TbCoreDeviceRpcService tbCoreDeviceRpcService) { - this.mainConsumer = tbCoreQueueProvider.getToCoreMsgConsumer(); - this.nfConsumer = tbCoreQueueProvider.getToCoreNotificationsMsgConsumer(); - this.actorContext = actorContext; + super(actorContext, encodingService, + tbCoreQueueProvider.getToCoreMsgConsumer(), tbCoreQueueProvider.getToCoreNotificationsMsgConsumer()); this.stateService = stateService; this.localSubscriptionService = localSubscriptionService; this.subscriptionManagerService = subscriptionManagerService; - this.encodingService = encodingService; this.tbCoreDeviceRpcService = tbCoreDeviceRpcService; } @PostConstruct public void init() { - this.mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-core-consumer")); - this.notificationsConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-core-notifications-consumer")); + super.init("tb-core-consumer", "tb-core-notifications-consumer"); } @Override - public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { - if (partitionChangeEvent.getServiceKey().getServiceType() == ServiceType.TB_CORE) { - log.info("Subscribing to partitions: {}", partitionChangeEvent.getPartitions()); - this.mainConsumer.subscribe(partitionChangeEvent.getPartitions()); - } - } - - @EventListener(ApplicationReadyEvent.class) - public void onApplicationEvent(ApplicationReadyEvent event) { - log.info("Subscribing to notifications: {}", mainConsumer.getTopic()); - this.nfConsumer.subscribe(); - launchNotificationsConsumer(); - launchMainConsumer(); - } - - private void launchMainConsumer() { + protected void launchMainConsumer() { + log.info("Launching main consumer"); mainConsumerExecutor.execute(() -> { while (!stopped) { try { @@ -134,6 +112,7 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { ConcurrentMap> failedMap = new ConcurrentHashMap<>(); CountDownLatch processingTimeoutLatch = new CountDownLatch(1); pendingMap.forEach((id, msg) -> { + log.info("[{}] Creating main callback for message: {}", id, msg.getValue()); TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>(), failedMap); try { ToCoreMsg toCoreMsg = msg.getValue(); @@ -177,57 +156,38 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { }); } - private void launchNotificationsConsumer() { - notificationsConsumerExecutor.execute(() -> { - while (!stopped) { - try { - List> msgs = nfConsumer.poll(pollDuration); - if (msgs.isEmpty()) { - continue; - } - ConcurrentMap> pendingMap = msgs.stream().collect( - Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity())); - ConcurrentMap> failedMap = new ConcurrentHashMap<>(); - CountDownLatch processingTimeoutLatch = new CountDownLatch(1); - pendingMap.forEach((id, msg) -> { - TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>(), failedMap); - try { - ToCoreNotificationMsg toCoreMsg = msg.getValue(); - if (toCoreMsg.hasToLocalSubscriptionServiceMsg()) { - log.trace("[{}] Forwarding message to local subscription service {}", id, toCoreMsg.getToLocalSubscriptionServiceMsg()); - forwardToLocalSubMgrService(toCoreMsg.getToLocalSubscriptionServiceMsg(), callback); - } else if (toCoreMsg.hasFromDeviceRpcResponse()) { - log.trace("[{}] Forwarding message to RPC service {}", id, toCoreMsg.getFromDeviceRpcResponse()); - forwardToCoreRpcService(toCoreMsg.getFromDeviceRpcResponse(), callback); - } else if (toCoreMsg.getComponentLifecycleMsg() != null && !toCoreMsg.getComponentLifecycleMsg().isEmpty()) { - Optional actorMsg = encodingService.decode(toCoreMsg.getComponentLifecycleMsg().toByteArray()); - if (actorMsg.isPresent()) { - log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); - actorContext.getAppActor().tell(actorMsg.get(), ActorRef.noSender()); - } - callback.onSuccess(); - } - } catch (Throwable e) { - log.warn("[{}] Failed to process notification: {}", id, msg, e); - callback.onFailure(e); - } - }); - if (!processingTimeoutLatch.await(packProcessingTimeout, TimeUnit.MILLISECONDS)) { - pendingMap.forEach((id, msg) -> log.warn("[{}] Timeout to process notification: {}", id, msg.getValue())); - failedMap.forEach((id, msg) -> log.warn("[{}] Failed to process notification: {}", id, msg.getValue())); - } - nfConsumer.commit(); - } catch (Exception e) { - log.warn("Failed to obtain notifications from queue.", e); - try { - Thread.sleep(pollDuration); - } catch (InterruptedException e2) { - log.trace("Failed to wait until the server has capacity to handle new notifications", e2); - } - } + @Override + protected ServiceType getServiceType() { + return ServiceType.TB_CORE; + } + + @Override + protected long getNotificationPollDuration() { + return pollDuration; + } + + @Override + protected long getNotificationPackProcessingTimeout() { + return packProcessingTimeout; + } + + @Override + protected void handleNotification(UUID id, TbProtoQueueMsg msg, TbMsgCallback callback) throws Exception { + ToCoreNotificationMsg toCoreMsg = msg.getValue(); + if (toCoreMsg.hasToLocalSubscriptionServiceMsg()) { + log.trace("[{}] Forwarding message to local subscription service {}", id, toCoreMsg.getToLocalSubscriptionServiceMsg()); + forwardToLocalSubMgrService(toCoreMsg.getToLocalSubscriptionServiceMsg(), callback); + } else if (toCoreMsg.hasFromDeviceRpcResponse()) { + log.trace("[{}] Forwarding message to RPC service {}", id, toCoreMsg.getFromDeviceRpcResponse()); + forwardToCoreRpcService(toCoreMsg.getFromDeviceRpcResponse(), callback); + } else if (toCoreMsg.getComponentLifecycleMsg() != null && !toCoreMsg.getComponentLifecycleMsg().isEmpty()) { + Optional actorMsg = encodingService.decode(toCoreMsg.getComponentLifecycleMsg().toByteArray()); + if (actorMsg.isPresent()) { + log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); + actorContext.getAppActor().tell(actorMsg.get(), ActorRef.noSender()); } - log.info("Tb Core Notifications Consumer stopped."); - }); + callback.onSuccess(); + } } private void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbMsgCallback callback) { @@ -245,23 +205,6 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { } } - @PreDestroy - public void destroy() { - stopped = true; - if (mainConsumer != null) { - mainConsumer.unsubscribe(); - } - if (nfConsumer != null) { - nfConsumer.unsubscribe(); - } - if (mainConsumerExecutor != null) { - mainConsumerExecutor.shutdownNow(); - } - if (notificationsConsumerExecutor != null) { - notificationsConsumerExecutor.shutdownNow(); - } - } - private void forwardToLocalSubMgrService(LocalSubscriptionServiceMsgProto msg, TbMsgCallback callback) { if (msg.hasSubUpdate()) { localSubscriptionService.onSubscriptionUpdate(msg.getSubUpdate().getSessionId(), TbSubscriptionUtils.fromProto(msg.getSubUpdate()), callback); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index b9d2b98d04..835da15ab7 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -19,49 +19,42 @@ import akka.actor.ActorRef; import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; -import org.springframework.boot.context.event.ApplicationReadyEvent; -import org.springframework.context.event.EventListener; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; -import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; +import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbMsgCallback; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.TbQueueConsumer; -import org.thingsboard.server.actors.ActorSystemContext; +import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; +import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.PartitionChangeEvent; -import org.thingsboard.server.common.msg.queue.ServiceType; -import org.thingsboard.server.gen.transport.TransportProtos.*; import org.thingsboard.server.queue.provider.TbRuleEngineQueueProvider; +import org.thingsboard.server.queue.util.TbMonolithOrRuleEngineComponent; import org.thingsboard.server.service.encoding.DataDecodingEncodingService; +import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingDecision; import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingResult; import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategy; import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory; import javax.annotation.PostConstruct; -import javax.annotation.PreDestroy; import java.util.List; import java.util.Optional; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.function.Function; import java.util.stream.Collectors; @Service -@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-rule-engine'") +@TbMonolithOrRuleEngineComponent @Slf4j -public class DefaultTbRuleEngineConsumerService implements TbRuleEngineConsumerService { +public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService implements TbRuleEngineConsumerService { @Value("${queue.rule_engine.poll_interval}") private long pollDuration; @@ -70,96 +63,24 @@ public class DefaultTbRuleEngineConsumerService implements TbRuleEngineConsumerS @Value("${queue.rule_engine.stats.enabled:false}") private boolean statsEnabled; - private final ActorSystemContext actorContext; - private final TbQueueConsumer> mainConsumer; - private final TbQueueConsumer> nfConsumer; private final TbCoreConsumerStats stats = new TbCoreConsumerStats(); private final TbRuleEngineProcessingStrategyFactory factory; - private final DataDecodingEncodingService encodingService; - private volatile ExecutorService mainConsumerExecutor; - private volatile ExecutorService notificationsConsumerExecutor; - private volatile boolean stopped = false; - public DefaultTbRuleEngineConsumerService(TbRuleEngineProcessingStrategyFactory factory, TbRuleEngineQueueProvider tbRuleEngineQueueProvider, ActorSystemContext actorContext, - DataDecodingEncodingService encodingService) { + public DefaultTbRuleEngineConsumerService(TbRuleEngineProcessingStrategyFactory factory, TbRuleEngineQueueProvider tbRuleEngineQueueProvider, + ActorSystemContext actorContext, DataDecodingEncodingService encodingService) { + super(actorContext, encodingService, + tbRuleEngineQueueProvider.getToRuleEngineMsgConsumer(), tbRuleEngineQueueProvider.getToRuleEngineNotificationsMsgConsumer()); this.factory = factory; - this.mainConsumer = tbRuleEngineQueueProvider.getToRuleEngineMsgConsumer(); - this.nfConsumer = tbRuleEngineQueueProvider.getToRuleEngineNotificationsMsgConsumer(); - this.actorContext = actorContext; - this.encodingService = encodingService; } @PostConstruct public void init() { - this.mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-rule-engine-consumer")); - this.notificationsConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-rule-engine-notifications-consumer")); + super.init("tb-rule-engine-consumer", "tb-rule-engine-notifications-consumer"); this.factory.newInstance(); } @Override - public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { - if (partitionChangeEvent.getServiceKey().getServiceType() == ServiceType.TB_RULE_ENGINE) { - log.info("Subscribing to partitions: {}", partitionChangeEvent.getPartitions()); - this.mainConsumer.subscribe(partitionChangeEvent.getPartitions()); - } - } - - @EventListener(ApplicationReadyEvent.class) - public void onApplicationEvent(ApplicationReadyEvent event) { - this.nfConsumer.subscribe(); - launchNotificationsConsumer(); - launchMainConsumer(); - } - - private void launchNotificationsConsumer() { - notificationsConsumerExecutor.execute(() -> { - while (!stopped) { - try { - List> msgs = nfConsumer.poll(pollDuration); - if (msgs.isEmpty()) { - continue; - } - ConcurrentMap> pendingMap = msgs.stream().collect( - Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity())); - ConcurrentMap> failedMap = new ConcurrentHashMap<>(); - CountDownLatch processingTimeoutLatch = new CountDownLatch(1); - pendingMap.forEach((id, msg) -> { - TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>(), failedMap); - try { - ToRuleEngineNotificationMsg toRuleEngineMsg = msg.getValue(); - if (toRuleEngineMsg.getComponentLifecycleMsg() != null && !toRuleEngineMsg.getComponentLifecycleMsg().isEmpty()) { - Optional actorMsg = encodingService.decode(toRuleEngineMsg.getComponentLifecycleMsg().toByteArray()); - if (actorMsg.isPresent()) { - log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); - actorContext.getAppActor().tell(actorMsg.get(), ActorRef.noSender()); - } - callback.onSuccess(); - } else { - callback.onSuccess(); - } - } catch (Throwable e) { - log.warn("[{}] Failed to process message: {}", id, msg, e); - callback.onFailure(e); - } - }); - if (!processingTimeoutLatch.await(packProcessingTimeout, TimeUnit.MILLISECONDS)) { - pendingMap.forEach((id, msg) -> log.warn("[{}] Timeout to process message: {}", id, msg.getValue())); - failedMap.forEach((id, msg) -> log.warn("[{}] Failed to process message: {}", id, msg.getValue())); - } - nfConsumer.commit(); - } catch (Exception e) { - log.warn("Failed to process messages from queue.", e); - try { - Thread.sleep(pollDuration); - } catch (InterruptedException e2) { - log.trace("Failed to wait until the server has capacity to handle new requests", e2); - } - } - } - }); - } - - private void launchMainConsumer() { + protected void launchMainConsumer() { mainConsumerExecutor.execute(() -> { while (!stopped) { try { @@ -184,6 +105,7 @@ public class DefaultTbRuleEngineConsumerService implements TbRuleEngineConsumerS CountDownLatch processingTimeoutLatch = new CountDownLatch(1); allMap.forEach((id, msg) -> { + log.info("[{}] Creating main callback for message: {}", id, msg.getValue()); TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, allMap, successMap, failedMap); try { ToRuleEngineMsg toRuleEngineMsg = msg.getValue(); @@ -217,6 +139,36 @@ public class DefaultTbRuleEngineConsumerService implements TbRuleEngineConsumerS }); } + @Override + protected ServiceType getServiceType() { + return ServiceType.TB_RULE_ENGINE; + } + + @Override + protected long getNotificationPollDuration() { + return pollDuration; + } + + @Override + protected long getNotificationPackProcessingTimeout() { + return packProcessingTimeout; + } + + @Override + protected void handleNotification(UUID id, TbProtoQueueMsg msg, TbMsgCallback callback) throws Exception { + ToRuleEngineNotificationMsg nfMsg = msg.getValue(); + if (nfMsg.getComponentLifecycleMsg() != null && !nfMsg.getComponentLifecycleMsg().isEmpty()) { + Optional actorMsg = encodingService.decode(nfMsg.getComponentLifecycleMsg().toByteArray()); + if (actorMsg.isPresent()) { + log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); + actorContext.getAppActor().tell(actorMsg.get(), ActorRef.noSender()); + } + callback.onSuccess(); + } else { + callback.onSuccess(); + } + } + private void forwardToRuleEngineActor(TenantId tenantId, ByteString tbMsgData, TbMsgCallback callback) { TbMsg tbMsg = TbMsg.fromBytes(tbMsgData.toByteArray(), callback); actorContext.getAppActor().tell(new QueueToRuleEngineMsg(tenantId, tbMsg), ActorRef.noSender()); @@ -233,20 +185,4 @@ public class DefaultTbRuleEngineConsumerService implements TbRuleEngineConsumerS } } - @PreDestroy - public void destroy() { - stopped = true; - if (mainConsumer != null) { - mainConsumer.unsubscribe(); - } - if (nfConsumer != null) { - nfConsumer.unsubscribe(); - } - if (mainConsumerExecutor != null) { - mainConsumerExecutor.shutdownNow(); - } - if (notificationsConsumerExecutor != null) { - notificationsConsumerExecutor.shutdownNow(); - } - } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java new file mode 100644 index 0000000000..e99770faf8 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java @@ -0,0 +1,150 @@ +/** + * Copyright © 2016-2020 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.queue.processing; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.ApplicationListener; +import org.springframework.context.event.EventListener; +import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.actors.ActorSystemContext; +import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TbMsgCallback; +import org.thingsboard.server.queue.TbQueueConsumer; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.discovery.PartitionChangeEvent; +import org.thingsboard.server.service.encoding.DataDecodingEncodingService; +import org.thingsboard.server.service.queue.MsgPackCallback; + +import javax.annotation.PreDestroy; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.function.Function; +import java.util.stream.Collectors; + +@Slf4j +public abstract class AbstractConsumerService implements ApplicationListener { + + protected volatile ExecutorService mainConsumerExecutor; + private volatile ExecutorService notificationsConsumerExecutor; + protected volatile boolean stopped = false; + + protected final ActorSystemContext actorContext; + protected final DataDecodingEncodingService encodingService; + protected final TbQueueConsumer> mainConsumer; + protected final TbQueueConsumer> nfConsumer; + + public AbstractConsumerService(ActorSystemContext actorContext, DataDecodingEncodingService encodingService, TbQueueConsumer> mainConsumer, TbQueueConsumer> nfConsumer) { + this.actorContext = actorContext; + this.encodingService = encodingService; + this.mainConsumer = mainConsumer; + this.nfConsumer = nfConsumer; + } + + public void init(String mainConsumerThreadName, String nfConsumerThreadName) { + this.mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName(mainConsumerThreadName)); + this.notificationsConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName(nfConsumerThreadName)); + } + + @Override + public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { + if (partitionChangeEvent.getServiceKey().getServiceType() == getServiceType()) { + log.info("Subscribing to partitions: {}", partitionChangeEvent.getPartitions()); + this.mainConsumer.subscribe(partitionChangeEvent.getPartitions()); + } + } + + @EventListener(ApplicationReadyEvent.class) + public void onApplicationEvent(ApplicationReadyEvent event) { + log.info("Subscribing to notifications: {}", nfConsumer.getTopic()); + this.nfConsumer.subscribe(); + launchNotificationsConsumer(); + launchMainConsumer(); + } + + protected abstract ServiceType getServiceType(); + + protected abstract void launchMainConsumer(); + + protected abstract long getNotificationPollDuration(); + + protected abstract long getNotificationPackProcessingTimeout(); + + protected void launchNotificationsConsumer() { + notificationsConsumerExecutor.execute(() -> { + while (!stopped) { + try { + List> msgs = nfConsumer.poll(getNotificationPollDuration()); + if (msgs.isEmpty()) { + continue; + } + ConcurrentMap> pendingMap = msgs.stream().collect( + Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity())); + ConcurrentMap> failedMap = new ConcurrentHashMap<>(); + CountDownLatch processingTimeoutLatch = new CountDownLatch(1); + pendingMap.forEach((id, msg) -> { + log.info("[{}] Creating notification callback for message: {}", id, msg.getValue()); + TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>(), failedMap); + try { + handleNotification(id, msg, callback); + } catch (Throwable e) { + log.warn("[{}] Failed to process notification: {}", id, msg, e); + callback.onFailure(e); + } + }); + if (!processingTimeoutLatch.await(getNotificationPackProcessingTimeout(), TimeUnit.MILLISECONDS)) { + pendingMap.forEach((id, msg) -> log.warn("[{}] Timeout to process notification: {}", id, msg.getValue())); + failedMap.forEach((id, msg) -> log.warn("[{}] Failed to process notification: {}", id, msg.getValue())); + } + nfConsumer.commit(); + } catch (Exception e) { + log.warn("Failed to obtain notifications from queue.", e); + try { + Thread.sleep(getNotificationPollDuration()); + } catch (InterruptedException e2) { + log.trace("Failed to wait until the server has capacity to handle new notifications", e2); + } + } + } + log.info("Tb Core Notifications Consumer stopped."); + }); + } + + protected abstract void handleNotification(UUID id, TbProtoQueueMsg msg, TbMsgCallback callback) throws Exception; + + @PreDestroy + public void destroy() { + stopped = true; + if (mainConsumer != null) { + mainConsumer.unsubscribe(); + } + if (nfConsumer != null) { + nfConsumer.unsubscribe(); + } + if (mainConsumerExecutor != null) { + mainConsumerExecutor.shutdownNow(); + } + if (notificationsConsumerExecutor != null) { + notificationsConsumerExecutor.shutdownNow(); + } + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java index ed55ee11df..cc9cc4dd10 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java @@ -36,6 +36,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.queue.TbClusterService; import javax.annotation.PostConstruct; @@ -54,7 +55,7 @@ import java.util.function.Consumer; */ @Service @Slf4j -@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core'") +@TbMonolithOrCoreComponent public class DefaultTbCoreDeviceRpcService implements TbCoreDeviceRpcService { private static final ObjectMapper json = new ObjectMapper(); @@ -79,7 +80,7 @@ public class DefaultTbCoreDeviceRpcService implements TbCoreDeviceRpcService { this.actorContext = actorContext; } - @Autowired + @Autowired(required = false) public void setTbRuleEngineRpcService(Optional tbRuleEngineRpcService) { this.tbRuleEngineRpcService = tbRuleEngineRpcService; } diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index 6481f1c2b6..f2390b9629 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -17,8 +17,6 @@ package org.thingsboard.server.service.rpc; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; -import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.rule.engine.api.RpcError; @@ -31,6 +29,7 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.util.TbMonolithOrRuleEngineComponent; import org.thingsboard.server.service.queue.TbClusterService; import javax.annotation.PostConstruct; @@ -45,7 +44,7 @@ import java.util.concurrent.TimeUnit; import java.util.function.Consumer; @Service -@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-rule-engine'") +@TbMonolithOrRuleEngineComponent @Slf4j public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcService { @@ -67,7 +66,7 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi this.serviceInfoProvider = serviceInfoProvider; } - @Autowired + @Autowired(required = false) public void setTbCoreRpcService(Optional tbCoreRpcService) { this.tbCoreRpcService = tbCoreRpcService; } diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAuthenticationProvider.java b/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAuthenticationProvider.java index e359befb30..ac5b3625f5 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAuthenticationProvider.java +++ b/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAuthenticationProvider.java @@ -39,6 +39,7 @@ import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.dao.audit.AuditLogService; import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.user.UserService; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.model.UserPrincipal; import org.thingsboard.server.service.security.system.SystemSecurityService; @@ -46,6 +47,7 @@ import ua_parser.Client; import java.util.UUID; + @Component @Slf4j public class RestAuthenticationProvider implements AuthenticationProvider { diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationFailureHandler.java b/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationFailureHandler.java index 726ee76a2d..406b078e6f 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationFailureHandler.java +++ b/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationFailureHandler.java @@ -20,6 +20,7 @@ import org.springframework.security.core.AuthenticationException; import org.springframework.security.web.authentication.AuthenticationFailureHandler; import org.springframework.stereotype.Component; import org.thingsboard.server.exception.ThingsboardErrorResponseHandler; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import javax.servlet.ServletException; import javax.servlet.http.HttpServletRequest; diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationSuccessHandler.java b/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationSuccessHandler.java index aa55818084..80cc8ec6e0 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationSuccessHandler.java +++ b/application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAwareAuthenticationSuccessHandler.java @@ -23,6 +23,7 @@ import org.springframework.security.core.Authentication; import org.springframework.security.web.WebAttributes; import org.springframework.security.web.authentication.AuthenticationSuccessHandler; import org.springframework.stereotype.Component; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.security.auth.jwt.RefreshTokenRepository; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.model.token.JwtToken; diff --git a/application/src/main/java/org/thingsboard/server/service/security/device/DefaultDeviceAuthService.java b/application/src/main/java/org/thingsboard/server/service/security/device/DefaultDeviceAuthService.java index c58a664790..f686cc2037 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/device/DefaultDeviceAuthService.java +++ b/application/src/main/java/org/thingsboard/server/service/security/device/DefaultDeviceAuthService.java @@ -24,16 +24,21 @@ import org.thingsboard.server.common.transport.auth.DeviceAuthResult; import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; @Service +@TbMonolithOrCoreComponent @Slf4j public class DefaultDeviceAuthService implements DeviceAuthService { - @Autowired - DeviceService deviceService; + private final DeviceService deviceService; - @Autowired - DeviceCredentialsService deviceCredentialsService; + private final DeviceCredentialsService deviceCredentialsService; + + public DefaultDeviceAuthService(DeviceService deviceService, DeviceCredentialsService deviceCredentialsService) { + this.deviceService = deviceService; + this.deviceCredentialsService = deviceCredentialsService; + } @Override public DeviceAuthResult process(DeviceCredentialsFilter credentialsFilter) { diff --git a/application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java b/application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java index ff22c00154..f8907b3188 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java +++ b/application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java @@ -28,6 +28,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.config.JwtSettings; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.model.UserPrincipal; diff --git a/application/src/main/java/org/thingsboard/server/service/security/permission/AccessControlService.java b/application/src/main/java/org/thingsboard/server/service/security/permission/AccessControlService.java index 40d19da764..5412ed6def 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/permission/AccessControlService.java +++ b/application/src/main/java/org/thingsboard/server/service/security/permission/AccessControlService.java @@ -15,11 +15,9 @@ */ package org.thingsboard.server.service.security.permission; -import org.thingsboard.server.common.data.HasCustomerId; import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.service.security.model.SecurityUser; public interface AccessControlService { diff --git a/application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java b/application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java index 9777af7b8e..be99cc1aae 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java +++ b/application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java @@ -30,6 +30,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.dao.customer.CustomerService; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.security.model.SecurityUser; import java.util.*; diff --git a/application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java b/application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java index 11857f25b6..ca6d1f531d 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java +++ b/application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java @@ -46,6 +46,7 @@ import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.settings.AdminSettingsService; import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.dao.user.UserServiceImpl; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.security.exception.UserPasswordExpiredException; import org.thingsboard.server.common.data.security.model.SecuritySettings; import org.thingsboard.server.common.data.security.model.UserPasswordPolicy; diff --git a/application/src/main/java/org/thingsboard/server/service/session/DefaultDeviceSessionCacheService.java b/application/src/main/java/org/thingsboard/server/service/session/DefaultDeviceSessionCacheService.java index 952c63c9b9..35c940ef86 100644 --- a/application/src/main/java/org/thingsboard/server/service/session/DefaultDeviceSessionCacheService.java +++ b/application/src/main/java/org/thingsboard/server/service/session/DefaultDeviceSessionCacheService.java @@ -21,6 +21,7 @@ import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import java.util.Collections; import java.util.UUID; @@ -31,6 +32,7 @@ import static org.thingsboard.server.common.data.CacheConstants.SESSIONS_CACHE; * Created by ashvayka on 29.10.18. */ @Service +@TbMonolithOrCoreComponent @Slf4j public class DefaultDeviceSessionCacheService implements DeviceSessionCacheService { diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 91a49aae1e..986c2cfae7 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -60,6 +60,7 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.common.msg.queue.TbMsgCallback; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.annotation.Nullable; @@ -89,6 +90,7 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; * Created by ashvayka on 01.05.18. */ @Service +@TbMonolithOrCoreComponent @Slf4j public class DefaultDeviceStateService implements DeviceStateService { @@ -109,7 +111,8 @@ public class DefaultDeviceStateService implements DeviceStateService { private final TimeseriesService tsService; private final TbQueueProducerProvider producerProvider; private final PartitionService partitionService; - private final TelemetrySubscriptionService tsSubService; + + private TelemetrySubscriptionService tsSubService; @Value("${state.defaultInactivityTimeoutInSec}") @Getter @@ -137,13 +140,17 @@ public class DefaultDeviceStateService implements DeviceStateService { public DefaultDeviceStateService(TenantService tenantService, DeviceService deviceService, AttributesService attributesService, TimeseriesService tsService, - TbQueueProducerProvider producerProvider, PartitionService partitionService, TelemetrySubscriptionService tsSubService) { + TbQueueProducerProvider producerProvider, PartitionService partitionService) { this.tenantService = tenantService; this.deviceService = deviceService; this.attributesService = attributesService; this.tsService = tsService; this.producerProvider = producerProvider; this.partitionService = partitionService; + } + + @Autowired + public void setTsSubService(TelemetrySubscriptionService tsSubService) { this.tsSubService = tsSubService; } @@ -188,9 +195,8 @@ public class DefaultDeviceStateService implements DeviceStateService { } @Override - public void onDeviceActivity(DeviceId deviceId) { - deviceLastReportedActivity.put(deviceId, System.currentTimeMillis()); - long lastReportedActivity = deviceLastReportedActivity.getOrDefault(deviceId, 0L); + public void onDeviceActivity(DeviceId deviceId, long lastReportedActivity) { + deviceLastReportedActivity.put(deviceId, lastReportedActivity); long lastSavedActivity = deviceLastSavedActivity.getOrDefault(deviceId, 0L); if (lastReportedActivity > 0 && lastReportedActivity > lastSavedActivity) { DeviceStateData stateData = getOrFetchDeviceStateData(deviceId); @@ -494,7 +500,7 @@ public class DefaultDeviceStateService implements DeviceStateService { try { TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), msgType, stateData.getDeviceId(), stateData.getMetaData().copy(), TbMsgDataType.JSON , json.writeValueAsString(state) - , null, null, null); + , null, null, TbMsgCallback.EMPTY); TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, stateData.getTenantId(), stateData.getDeviceId()); TransportProtos.ToRuleEngineMsg msg = TransportProtos.ToRuleEngineMsg.newBuilder() .setTenantIdMSB(stateData.getTenantId().getId().getMostSignificantBits()) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java index 430dd77a36..c7aadda963 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java @@ -35,7 +35,7 @@ public interface DeviceStateService extends ApplicationListener currentPartitions = ConcurrentHashMap.newKeySet(); - @Autowired - private AttributesService attrService; - - @Autowired - private TimeseriesService tsService; - - @Autowired - private TbQueueProducerProvider producerProvider; - - @Autowired - private PartitionService partitionService; - - @Autowired - private SubscriptionManagerService subscriptionManagerService; + private final AttributesService attrService; + private final TimeseriesService tsService; + private final TbQueueProducerProvider producerProvider; + private final PartitionService partitionService; + private Optional subscriptionManagerService; private TbQueueProducer> toCoreProducer; private ExecutorService tsCallBackExecutor; private ExecutorService wsCallBackExecutor; + public DefaultTelemetrySubscriptionService(AttributesService attrService, + TimeseriesService tsService, + TbQueueProducerProvider producerProvider, + PartitionService partitionService) { + this.attrService = attrService; + this.tsService = tsService; + this.producerProvider = producerProvider; + this.partitionService = partitionService; + } + + @Autowired(required = false) + public void setSubscriptionManagerService(Optional subscriptionManagerService) { + this.subscriptionManagerService = subscriptionManagerService; + } + @PostConstruct public void initExecutor() { tsCallBackExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("ts-service-ts-callback")); @@ -158,7 +165,11 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio private void onAttributesUpdate(TenantId tenantId, EntityId entityId, String scope, List attributes) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); if (currentPartitions.contains(tpi)) { - subscriptionManagerService.onAttributesUpdate(tenantId, entityId, scope, attributes, TbMsgCallback.EMPTY); + if (subscriptionManagerService.isPresent()) { + subscriptionManagerService.get().onAttributesUpdate(tenantId, entityId, scope, attributes, TbMsgCallback.EMPTY); + } else { + log.warn("Possible misconfiguration because subscriptionManagerService is null!"); + } } else { TransportProtos.ToCoreMsg toCoreMsg = TbSubscriptionUtils.toAttributesUpdateProto(tenantId, entityId, scope, attributes); toCoreProducer.send(tpi, new TbProtoQueueMsg<>(entityId.getId(), toCoreMsg), null); @@ -168,7 +179,11 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio private void onTimeSeriesUpdate(TenantId tenantId, EntityId entityId, List ts) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); if (currentPartitions.contains(tpi)) { - subscriptionManagerService.onTimeSeriesUpdate(tenantId, entityId, ts, TbMsgCallback.EMPTY); + if (subscriptionManagerService.isPresent()) { + subscriptionManagerService.get().onTimeSeriesUpdate(tenantId, entityId, ts, TbMsgCallback.EMPTY); + } else { + log.warn("Possible misconfiguration because subscriptionManagerService is null!"); + } } else { TransportProtos.ToCoreMsg toCoreMsg = TbSubscriptionUtils.toTimeseriesUpdateProto(tenantId, entityId, ts); toCoreProducer.send(tpi, new TbProtoQueueMsg<>(entityId.getId(), toCoreMsg), null); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java index 876ea4f2eb..b4c15ee8cc 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java @@ -43,6 +43,7 @@ import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.util.TenantRateLimitException; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.security.AccessValidator; import org.thingsboard.server.service.security.ValidationCallback; import org.thingsboard.server.service.security.ValidationResult; @@ -83,6 +84,7 @@ import java.util.stream.Collectors; * Created by ashvayka on 27.03.18. */ @Service +@TbMonolithOrCoreComponent @Slf4j public class DefaultTelemetryWebSocketService implements TelemetryWebSocketService { diff --git a/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java b/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java index b9424d49d9..b402f88896 100644 --- a/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java +++ b/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java @@ -24,6 +24,8 @@ import org.thingsboard.rule.engine.api.RuleChainTransactionService; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.cluster.ServerAddress; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; +import org.thingsboard.server.queue.util.TbMonolithOrRuleEngineComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import javax.annotation.PostConstruct; @@ -44,11 +46,11 @@ import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; @Service +@TbMonolithOrRuleEngineComponent @Slf4j public class BaseRuleChainTransactionService implements RuleChainTransactionService { - @Autowired - private DbCallbackExecutorService callbackExecutor; + private final DbCallbackExecutorService callbackExecutor; @Value("${actors.rule.transaction.queue_size}") private int finalQueueSize; @@ -61,6 +63,10 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ private ExecutorService timeoutExecutor; + public BaseRuleChainTransactionService(DbCallbackExecutorService callbackExecutor) { + this.callbackExecutor = callbackExecutor; + } + @PostConstruct public void init() { timeoutExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("rule-chain-transaction")); diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java index 52d748b96b..ab257f00f9 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java @@ -27,6 +27,7 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos.DeviceActorToTransportMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import java.util.UUID; import java.util.function.Consumer; @@ -35,7 +36,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; @Slf4j @Service -@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core'") +@TbMonolithOrCoreComponent public class DefaultTbCoreToTransportService implements TbCoreToTransportService { private final TbQueueProducer> tbTransportProducer; diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 42ee4c4619..dbff5bda86 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -41,6 +41,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponse import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceTokenRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.state.DeviceStateService; @@ -52,6 +53,7 @@ import java.util.concurrent.locks.ReentrantLock; */ @Slf4j @Service +@TbMonolithOrCoreComponent public class DefaultTransportApiService implements TransportApiService { private static final ObjectMapper mapper = new ObjectMapper(); diff --git a/application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/TbCoreTransportApiService.java similarity index 93% rename from application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java rename to application/src/main/java/org/thingsboard/server/service/transport/TbCoreTransportApiService.java index 4c9c801e83..446f7ba81a 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/TbCoreTransportApiService.java @@ -28,6 +28,8 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg; import org.thingsboard.server.queue.provider.TbCoreQueueProvider; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @@ -38,9 +40,8 @@ import java.util.concurrent.*; */ @Slf4j @Service -//TODO 2.5: This Confitional annotation should be removed, and Service renamed to something meaningful -//@ConditionalOnProperty(prefix = "transport", value = "type", havingValue = "remote") -public class RemoteTransportApiService { +@TbMonolithOrCoreComponent +public class TbCoreTransportApiService { private final TbCoreQueueProvider tbCoreQueueProvider; private final TransportApiService transportApiService; @@ -58,7 +59,7 @@ public class RemoteTransportApiService { private TbQueueResponseTemplate, TbProtoQueueMsg> transportApiTemplate; - public RemoteTransportApiService(TbCoreQueueProvider tbCoreQueueProvider, TransportApiService transportApiService) { + public TbCoreTransportApiService(TbCoreQueueProvider tbCoreQueueProvider, TransportApiService transportApiService) { this.tbCoreQueueProvider = tbCoreQueueProvider; this.transportApiService = transportApiService; } diff --git a/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java b/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java index 1b8173c424..8844c1fe5b 100644 --- a/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java @@ -24,6 +24,7 @@ import org.springframework.stereotype.Service; import org.springframework.web.client.RestTemplate; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.UpdateMessage; +import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @@ -38,6 +39,7 @@ import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @Service +@TbMonolithOrCoreComponent @Slf4j public class DefaultUpdateService implements UpdateService { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 7a80a26713..dd9e6ad3ec 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -205,8 +205,6 @@ sql: # Actor system parameters actors: - cluster: - grpc_callback_thread_pool_size: "${ACTORS_CLUSTER_GRPC_CALLBACK_THREAD_POOL_SIZE:10}" tenant: create_components_on_init: "${ACTORS_TENANT_CREATE_COMPONENTS_ON_INIT:true}" session: @@ -370,7 +368,7 @@ spring: username: "${SPRING_DATASOURCE_USERNAME:postgres}" password: "${SPRING_DATASOURCE_PASSWORD:postgres}" hikari: - maximumPoolSize: "${SPRING_DATASOURCE_MAXIMUM_POOL_SIZE:50}" + maximumPoolSize: "${SPRING_DATASOURCE_MAXIMUM_POOL_SIZE:5}" # Audit log parameters audit-log: @@ -540,8 +538,6 @@ queue: response_poll_interval: "${TB_QUEUE_TRANSPORT_RESPONSE_POLL_INTERVAL_MS:25}" core: topic: "${TB_QUEUE_CORE_TOPIC:tb.core}" - # For high priority notifications that require minimum latency and processing time - notifications_topic: "${TB_QUEUE_CORE_NOTIFICATIONS_TOPIC:tb.rule-engine.notifications}" poll_interval: "${TB_QUEUE_CORE_POLL_INTERVAL_MS:25}" partitions: "${TB_QUEUE_CORE_PARTITIONS:10}" pack_processing_timeout: "${TB_QUEUE_CORE_PACK_PROCESSING_TIMEOUT_MS:60000}" @@ -550,8 +546,6 @@ queue: print_interval_ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:10000}" rule_engine: topic: "${TB_QUEUE_RULE_ENGINE_TOPIC:tb.rule-engine}" - # For high priority notifications that require minimum latency and processing time - notifications_topic: "${TB_QUEUE_RULE_ENGINE_NOTIFICATIONS_TOPIC:tb.rule-engine.notifications}" poll_interval: "${TB_QUEUE_RULE_ENGINE_POLL_INTERVAL_MS:25}" partitions: "${TB_QUEUE_RULE_ENGINE_PARTITIONS:10}" pack_processing_timeout: "${TB_QUEUE_RULE_ENGINE_PACK_PROCESSING_TIMEOUT_MS:60000}" @@ -573,4 +567,5 @@ service: type: "${TB_SERVICE_TYPE:monolith}" # monolith or tb-core or tb-rule-engine or tb-transport # Unique id for this service (autogenerated if empty) id: "${TB_SERVICE_ID:}" - tenant_id: "${TB_SERVICE_TENANT_ID:}" # empty or specific tenant id. \ No newline at end of file + tenant_id: "${TB_SERVICE_TENANT_ID:}" # empty or specific tenant id. + diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceKey.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceKey.java index b68717730b..aa6eb27323 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceKey.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceKey.java @@ -16,10 +16,12 @@ package org.thingsboard.server.common.msg.queue; import lombok.Getter; +import lombok.ToString; import org.thingsboard.server.common.data.id.TenantId; import java.util.Objects; +@ToString public class ServiceKey { @Getter private final ServiceType serviceType; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java index 02ad1d0528..bf0fffe404 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java @@ -17,11 +17,13 @@ package org.thingsboard.server.common.msg.queue; import lombok.Builder; import lombok.Getter; +import lombok.ToString; import org.thingsboard.server.common.data.id.TenantId; import java.util.Objects; import java.util.Optional; +@ToString public class TopicPartitionInfo { private final String topic; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueProvider.java index 56cf4fb55f..3566b29478 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueProvider.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.queue.provider; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -37,9 +36,12 @@ import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.kafka.TBKafkaConsumerTemplate; import org.thingsboard.server.queue.kafka.TBKafkaProducerTemplate; import org.thingsboard.server.queue.kafka.TbKafkaSettings; +import org.thingsboard.server.queue.util.TbKafkaQueue; +import org.thingsboard.server.queue.util.TbMonolithComponent; @Component -@ConditionalOnExpression("'${queue.type:null}'=='kafka' && '${service.type:null}'=='monolith'") +@TbMonolithComponent +@TbKafkaQueue public class KafkaMonolithQueueProvider implements TbCoreQueueProvider, TbRuleEngineQueueProvider { private final PartitionService partitionService; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java index 177f5f3c3d..12aa66838e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.queue.provider; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -37,9 +36,12 @@ import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.kafka.TBKafkaConsumerTemplate; import org.thingsboard.server.queue.kafka.TBKafkaProducerTemplate; import org.thingsboard.server.queue.kafka.TbKafkaSettings; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.queue.util.TbKafkaQueue; @Component -@ConditionalOnExpression("'${queue.type:null}'=='kafka' && '${service.type:null}'=='tb-core'") +@TbCoreComponent +@TbKafkaQueue public class KafkaTbCoreQueueProvider implements TbCoreQueueProvider { private final PartitionService partitionService; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java index 36ffd05b48..f295f61d6a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.queue.provider; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -35,9 +34,12 @@ import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.kafka.TBKafkaConsumerTemplate; import org.thingsboard.server.queue.kafka.TBKafkaProducerTemplate; import org.thingsboard.server.queue.kafka.TbKafkaSettings; +import org.thingsboard.server.queue.util.TbKafkaQueue; +import org.thingsboard.server.queue.util.TbRuleEngineComponent; @Component -@ConditionalOnExpression("'${queue.type:null}'=='kafka' && '${service.type:null}'=='tb-rule-engine'") +@TbKafkaQueue +@TbRuleEngineComponent public class KafkaTbRuleEngineQueueProvider implements TbRuleEngineQueueProvider { private final PartitionService partitionService; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProvider.java index bbbef16f65..e74c07152b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProvider.java @@ -15,7 +15,7 @@ */ package org.thingsboard.server.queue.provider; -import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.gen.transport.TransportProtos.*; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; @@ -50,7 +50,7 @@ public interface TbCoreQueueProvider { * * @return */ - TbQueueProducer> getRuleEngineNotificationsMsgProducer(); + TbQueueProducer> getRuleEngineNotificationsMsgProducer(); /** * Used to push messages to other instances of TB Core Service @@ -64,7 +64,7 @@ public interface TbCoreQueueProvider { * * @return */ - TbQueueProducer> getTbCoreNotificationsMsgProducer(); + TbQueueProducer> getTbCoreNotificationsMsgProducer(); /** * Used to consume messages by TB Core Service @@ -78,7 +78,7 @@ public interface TbCoreQueueProvider { * * @return */ - TbQueueConsumer> getToCoreNotificationsMsgConsumer(); + TbQueueConsumer> getToCoreNotificationsMsgConsumer(); /** * Used to consume Transport API Calls diff --git a/application/src/main/java/org/thingsboard/server/service/transport/ToTransportMsgEncoder.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbCoreComponent.java similarity index 60% rename from application/src/main/java/org/thingsboard/server/service/transport/ToTransportMsgEncoder.java rename to common/queue/src/main/java/org/thingsboard/server/queue/util/TbCoreComponent.java index 50b940ae48..244aae5394 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/ToTransportMsgEncoder.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbCoreComponent.java @@ -13,17 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.transport; +package org.thingsboard.server.queue.util; -import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; -import org.thingsboard.server.queue.kafka.TbKafkaEncoder; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; -/** - * Created by ashvayka on 05.10.18. - */ -public class ToTransportMsgEncoder implements TbKafkaEncoder { - @Override - public byte[] encode(ToTransportMsg value) { - return value.toByteArray(); - } +@ConditionalOnExpression("'${service.type:null}'=='tb-core'") +public @interface TbCoreComponent { } diff --git a/application/src/main/java/org/thingsboard/server/service/transport/ToRuleEngineMsgDecoder.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbKafkaQueue.java similarity index 53% rename from application/src/main/java/org/thingsboard/server/service/transport/ToRuleEngineMsgDecoder.java rename to common/queue/src/main/java/org/thingsboard/server/queue/util/TbKafkaQueue.java index 71739c0439..8d4af73492 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/ToRuleEngineMsgDecoder.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbKafkaQueue.java @@ -13,21 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.transport; +package org.thingsboard.server.queue.util; -import org.thingsboard.server.queue.TbQueueMsg; -import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; -import org.thingsboard.server.queue.kafka.TbKafkaDecoder; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; -import java.io.IOException; - -/** - * Created by ashvayka on 05.10.18. - */ -public class ToRuleEngineMsgDecoder implements TbKafkaDecoder { - - @Override - public ToRuleEngineMsg decode(TbQueueMsg msg) throws IOException { - return ToRuleEngineMsg.parseFrom(msg.getData()); - } +@ConditionalOnExpression("'${queue.type:null}'=='kafka'") +public @interface TbKafkaQueue { } diff --git a/application/src/main/java/org/thingsboard/server/service/executors/ClusterRpcCallbackExecutorService.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithComponent.java similarity index 53% rename from application/src/main/java/org/thingsboard/server/service/executors/ClusterRpcCallbackExecutorService.java rename to common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithComponent.java index f4b14144fc..c6dd9ebfff 100644 --- a/application/src/main/java/org/thingsboard/server/service/executors/ClusterRpcCallbackExecutorService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithComponent.java @@ -13,21 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.executors; +package org.thingsboard.server.queue.util; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import org.thingsboard.common.util.AbstractListeningExecutor; - -@Component -public class ClusterRpcCallbackExecutorService extends AbstractListeningExecutor { - - @Value("${actors.cluster.grpc_callback_thread_pool_size}") - private int grpcCallbackExecutorThreadPoolSize; - - @Override - protected int getThreadPollSize() { - return grpcCallbackExecutorThreadPoolSize; - } +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; +@ConditionalOnExpression("'${service.type:null}'=='monolith'") +public @interface TbMonolithComponent { } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrCoreComponent.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrCoreComponent.java new file mode 100644 index 0000000000..a0165e3386 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrCoreComponent.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2020 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.queue.util; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; + +@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core'") +public @interface TbMonolithOrCoreComponent { +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrRuleEngineComponent.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrRuleEngineComponent.java new file mode 100644 index 0000000000..83e8a2d87a --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbMonolithOrRuleEngineComponent.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2020 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.queue.util; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; + +@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-rule-engine'") +public @interface TbMonolithOrRuleEngineComponent { +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/TbRuleEngineComponent.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbRuleEngineComponent.java new file mode 100644 index 0000000000..9e79232906 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbRuleEngineComponent.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2020 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.queue.util; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; + +@ConditionalOnExpression("'${service.type:null}'=='tb-rule-engine'") +public @interface TbRuleEngineComponent { +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 0d96f94252..fb509e91af 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -410,13 +410,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void processDisconnect(ChannelHandlerContext ctx) { ctx.close(); log.info("[{}] Client disconnected!", sessionId); - if (deviceSessionCtx.isConnected()) { - transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.CLOSED), null); - transportService.deregisterSession(sessionInfo); - if (gatewaySessionHandler != null) { - gatewaySessionHandler.onGatewayDisconnect(); - } - } + doDisconnect(); } private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode) { @@ -485,9 +479,17 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void operationComplete(Future future) throws Exception { + doDisconnect(); + } + + private void doDisconnect() { if (deviceSessionCtx.isConnected()) { transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.CLOSED), null); transportService.deregisterSession(sessionInfo); + if (gatewaySessionHandler != null) { + gatewaySessionHandler.onGatewayDisconnect(); + } + deviceSessionCtx.setDisconnected(); } } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 9b8c23bf8b..3fdf19c741 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -45,6 +45,7 @@ import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.provider.TbTransportQueueProvider; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -92,6 +93,7 @@ public class DefaultTransportService implements TransportService { private final Gson gson = new Gson(); private final TbTransportQueueProvider queueProvider; + private final TbQueueProducerProvider producerProvider; private final PartitionService partitionService; protected TbQueueRequestTemplate, TbProtoQueueMsg> transportApiRequestTemplate; @@ -110,8 +112,9 @@ public class DefaultTransportService implements TransportService { private ExecutorService mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer")); private volatile boolean stopped = false; - public DefaultTransportService(TbTransportQueueProvider queueProvider, PartitionService partitionService) { + public DefaultTransportService(TbTransportQueueProvider queueProvider, TbQueueProducerProvider producerProvider, PartitionService partitionService) { this.queueProvider = queueProvider; + this.producerProvider = producerProvider; this.partitionService = partitionService; } @@ -126,8 +129,8 @@ public class DefaultTransportService implements TransportService { this.transportCallbackExecutor = Executors.newWorkStealingPool(20); this.schedulerExecutor.scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS); transportApiRequestTemplate = queueProvider.getTransportApiRequestTemplate(); - ruleEngineMsgProducer = queueProvider.getRuleEngineMsgProducer(); - tbCoreMsgProducer = queueProvider.getTbCoreMsgProducer(); + ruleEngineMsgProducer = producerProvider.getRuleEngineMsgProducer(); + tbCoreMsgProducer = producerProvider.getTbCoreMsgProducer(); transportNotificationsConsumer = queueProvider.getTransportNotificationsConsumer(); transportNotificationsConsumer.subscribe(); transportApiRequestTemplate.init(); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java index 7c590ec925..c312b3a4d9 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java @@ -35,6 +35,7 @@ public abstract class DeviceAwareSessionContext implements SessionContext { private volatile DeviceId deviceId; @Getter private volatile DeviceInfoProto deviceInfo; + private volatile boolean connected; public DeviceId getDeviceId() { return deviceId; @@ -42,10 +43,15 @@ public abstract class DeviceAwareSessionContext implements SessionContext { public void setDeviceInfo(DeviceInfoProto deviceInfo) { this.deviceInfo = deviceInfo; + this.connected = true; this.deviceId = new DeviceId(new UUID(deviceInfo.getDeviceIdMSB(), deviceInfo.getDeviceIdLSB())); } public boolean isConnected() { - return deviceInfo != null; + return connected; + } + + public void setDisconnected() { + this.connected = false; } }