diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java index 5c577dd39f..8b1dece58e 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java @@ -32,7 +32,6 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; -import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import javax.annotation.PostConstruct; @@ -53,11 +52,11 @@ import java.util.concurrent.atomic.AtomicInteger; @ConditionalOnProperty(name = "queue.core.housekeeper.enabled", havingValue = "true", matchIfMissing = true) public class HousekeeperReprocessingService { - private final DefaultHousekeeperService housekeeperService; + private final HousekeeperService housekeeperService; private final PartitionService partitionService; private final TbCoreQueueFactory queueFactory; - private final TbQueueProducer> producer; - private final TopicPartitionInfo submitTpi; + private TbQueueProducer> producer; + private TopicPartitionInfo submitTpi; @Value("${queue.core.housekeeper.reprocessing-start-delay-sec:300}") private int startDelay; @@ -74,18 +73,19 @@ public class HousekeeperReprocessingService { protected AtomicInteger cycle = new AtomicInteger(); private boolean stopped; - public HousekeeperReprocessingService(@Lazy DefaultHousekeeperService housekeeperService, - PartitionService partitionService, TbCoreQueueFactory queueFactory, - TbQueueProducerProvider producerProvider) { + public HousekeeperReprocessingService(@Lazy HousekeeperService housekeeperService, + PartitionService partitionService, + TbCoreQueueFactory queueFactory) { this.housekeeperService = housekeeperService; this.partitionService = partitionService; this.queueFactory = queueFactory; - this.producer = producerProvider.getHousekeeperReprocessingMsgProducer(); - this.submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); } @PostConstruct private void init() { + producer = queueFactory.createHousekeeperReprocessingMsgProducer(); + submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); + scheduler.scheduleWithFixedDelay(() -> { try { cycle.incrementAndGet(); diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java similarity index 78% rename from application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java rename to application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java index 589c570144..4c0c0e72d5 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java @@ -26,15 +26,11 @@ import org.thingsboard.server.common.data.housekeeper.HousekeeperTask; import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; import org.thingsboard.server.common.data.notification.rule.trigger.TaskProcessingFailureTrigger; import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; -import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.dao.housekeeper.HousekeeperService; import org.thingsboard.server.gen.transport.TransportProtos.HousekeeperTaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import org.thingsboard.server.queue.TbQueueConsumer; -import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; -import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.housekeeper.processor.HousekeeperTaskProcessor; @@ -43,6 +39,7 @@ import org.thingsboard.server.service.housekeeper.stats.HousekeeperStatsService; import javax.annotation.PreDestroy; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -55,16 +52,14 @@ import java.util.stream.Collectors; @Service @Slf4j @ConditionalOnProperty(name = "queue.core.housekeeper.enabled", havingValue = "true", matchIfMissing = true) -public class DefaultHousekeeperService implements HousekeeperService { +public class HousekeeperService { private final Map> taskProcessors; private final HousekeeperReprocessingService reprocessingService; - private final HousekeeperStatsService statsService; + private final Optional statsService; private final NotificationRuleProcessor notificationRuleProcessor; private final TbQueueConsumer> consumer; - private final TbQueueProducer> producer; - private final TopicPartitionInfo submitTpi; private final ExecutorService consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper-consumer")); private final ExecutorService executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper-task-processor")); @@ -76,18 +71,15 @@ public class DefaultHousekeeperService implements HousekeeperService { private boolean stopped; - public DefaultHousekeeperService(HousekeeperReprocessingService reprocessingService, - TbCoreQueueFactory queueFactory, - TbQueueProducerProvider producerProvider, - HousekeeperStatsService statsService, - NotificationRuleProcessor notificationRuleProcessor, - @Lazy List> taskProcessors) { + public HousekeeperService(HousekeeperReprocessingService reprocessingService, + TbCoreQueueFactory queueFactory, + Optional statsService, + NotificationRuleProcessor notificationRuleProcessor, + @Lazy List> taskProcessors) { this.reprocessingService = reprocessingService; this.statsService = statsService; this.notificationRuleProcessor = notificationRuleProcessor; this.consumer = queueFactory.createHousekeeperMsgConsumer(); - this.producer = producerProvider.getHousekeeperMsgProducer(); - this.submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); this.taskProcessors = taskProcessors.stream().collect(Collectors.toMap(HousekeeperTaskProcessor::getTaskType, p -> p)); } @@ -147,7 +139,7 @@ public class DefaultHousekeeperService implements HousekeeperService { }); future.get(taskProcessingTimeout, TimeUnit.MILLISECONDS); - statsService.reportProcessed(task.getTaskType(), msg); + statsService.ifPresent(statsService -> statsService.reportProcessed(task.getTaskType(), msg)); log.debug("[{}] Successfully {} task {}", task.getTenantId(), isNew(msg.getTask()) ? "processed" : "reprocessed", msg.getTask().getValue()); } catch (InterruptedException e) { throw e; @@ -163,7 +155,7 @@ public class DefaultHousekeeperService implements HousekeeperService { task.getTaskType(), msg.getTask().getAttempt(), task, error); reprocessingService.submitForReprocessing(msg, error); - statsService.reportFailure(task.getTaskType(), msg); + statsService.ifPresent(statsService -> statsService.reportFailure(task.getTaskType(), msg)); notificationRuleProcessor.process(TaskProcessingFailureTrigger.builder() .task(task) .error(error) @@ -172,22 +164,6 @@ public class DefaultHousekeeperService implements HousekeeperService { } } - @Override - public void submitTask(HousekeeperTask task) { - log.trace("[{}][{}][{}] Submitting task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType()); - /* - * using msg key as entity id so that msgs related to certain entity are pushed to same partition, - * e.g. on tenant deletion (entity id is tenant id), we need to clean up tenant entities in certain order - * */ - producer.send(submitTpi, new TbProtoQueueMsg<>(task.getEntityId().getId(), ToHousekeeperServiceMsg.newBuilder() - .setTask(HousekeeperTaskProto.newBuilder() - .setValue(JacksonUtil.toString(task)) - .setTs(task.getTs()) - .setAttempt(0) - .build()) - .build()), null); - } - private boolean isNew(HousekeeperTaskProto task) { return task.getErrorsCount() == 0; } diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java index 3f6e88a299..400bc36479 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java @@ -17,13 +17,14 @@ package org.thingsboard.server.service.housekeeper.stats; import lombok.Getter; import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; import org.thingsboard.server.common.stats.DefaultCounter; import org.thingsboard.server.common.stats.StatsCounter; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsType; -import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import java.util.ArrayList; @@ -31,11 +32,11 @@ import java.util.EnumMap; import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @Service @Slf4j +@ConditionalOnProperty(name = "queue.core.housekeeper.stats.enabled", havingValue = "true", matchIfMissing = true) public class HousekeeperStatsService { private final Map stats = new EnumMap<>(HousekeeperTaskType.class); @@ -46,7 +47,8 @@ public class HousekeeperStatsService { } } - @Scheduled(initialDelay = 60, fixedDelay = 60, timeUnit = TimeUnit.SECONDS) + @Scheduled(initialDelayString = "${queue.core.housekeeper.stats.print-interval-ms:60000}", + fixedDelayString = "${queue.core.housekeeper.stats.print-interval-ms:60000}") private void reportStats() { String statsStr = stats.values().stream().map(stats -> { String countersStr = stats.getCounters().stream() diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index a156762c71..5baa155ae9 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1401,9 +1401,11 @@ queue: # - key: max.poll.records # value: "${TB_QUEUE_KAFKA_SQ_MAX_POLL_RECORDS:1024}" tb_housekeeper: + # Amount of records to be returned in a single poll. For Housekeeper tasks topic, we should consume messages (tasks) one by one - key: max.poll.records value: "1" tb_housekeeper.reprocessing: + # Amount of records to be returned in a single poll. For Housekeeper reprocessing topic, we should consume messages (tasks) one by one - key: max.poll.records value: "1" other-inline: "${TB_QUEUE_KAFKA_OTHER_PROPERTIES:}" # In this section you can specify custom parameters (semicolon separated) for Kafka consumer/producer/admin # Example "metrics.recording.level:INFO;metrics.sample.window.ms:30000" @@ -1428,9 +1430,9 @@ queue: # Kafka properties for Version Control topic version-control: "${TB_QUEUE_KAFKA_VC_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}" # Kafka properties for Housekeeper tasks topic - housekeeper: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000;partitions:10;min.insync.replicas:1" - # Kafka properties for Housekeeper reprocessing topic; retention.ms is set to 90 days - housekeeper-reprocessing: "retention.ms:7776000000;segment.bytes:26214400;retention.bytes:1048576000;partitions:1;min.insync.replicas:1" + housekeeper: "${TB_QUEUE_KAFKA_HOUSEKEEPER_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000;partitions:10;min.insync.replicas:1}" + # Kafka properties for Housekeeper reprocessing topic; retention.ms is set to 90 days; partitions is set to 1 since only one reprocessing service is running at a time + housekeeper-reprocessing: "${TB_QUEUE_KAFKA_HOUSEKEEPER_REPROCESSING_TOPIC_PROPERTIES:retention.ms:7776000000;segment.bytes:26214400;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}" consumer-stats: # Prints lag between consumer group offset and last messages offset in Kafka topics enabled: "${TB_QUEUE_KAFKA_CONSUMER_STATS_ENABLED:true}" @@ -1587,14 +1589,27 @@ queue: # Statistics printing interval for Core microservices print-interval-ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:60000}" housekeeper: + # Enabled/disable Housekeeping service, which does cleanup of telemetry, attributes, events, etc. enabled: "${TB_HOUSEKEEPER_ENABLED:true}" + # Topic name for Housekeeper tasks topic: "${TB_HOUSEKEEPER_TOPIC:tb_housekeeper}" + # Topic name for Housekeeper tasks to be reprocessed reprocessing-topic: "${TB_HOUSEKEEPER_REPROCESSING_TOPIC:tb_housekeeper.reprocessing}" + # Poll interval for topics related to Housekeeper poll-interval-ms: "${TB_HOUSEKEEPER_POLL_INTERVAL_MS:500}" + # Timeout in milliseconds for task processing. Tasks that fail to finish on time will be submitted for reprocessing task-processing-timeout-ms: "${TB_HOUSEKEEPER_TASK_PROCESSING_TIMEOUT_MS:120000}" + # Delay in seconds for starting tasks reprocessing after start-up reprocessing-start-delay-sec: "${TB_HOUSEKEEPER_REPROCESSING_START_DELAY_SEC:300}" + # Period in seconds of tasks reprocessing task-reprocessing-delay-sec: "${TB_HOUSEKEEPER_TASK_REPROCESSING_DELAY_SEC:3600}" - max-reprocessing-attempts: "${TB_HOUSEKEEPER_MAX_REPROCESSING_ATTEMPTS:30}" + # Maximum amount of task reprocessing attempts. After exceeding, the task will be ignored until the next service start-up + max-reprocessing-attempts: "${TB_HOUSEKEEPER_MAX_REPROCESSING_ATTEMPTS:10}" + stats: + # Enable/disable statistics for Housekeeper + enabled: "${TB_HOUSEKEEPER_STATS_ENABLED:true}" + # Statistics printing interval for Housekeeper + print-interval-ms: "${TB_HOUSEKEEPER_STATS_PRINT_INTERVAL_MS:60000}" vc: # Default topic name for Kafka, RabbitMQ, etc. diff --git a/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java b/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java index ef3ce8c786..daf363467d 100644 --- a/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java @@ -56,6 +56,7 @@ import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.rule.RuleNode; +import org.thingsboard.server.common.msg.housekeeper.HousekeeperClient; import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.attributes.AttributesService; @@ -99,7 +100,9 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers. public class HousekeeperServiceTest extends AbstractControllerTest { @SpyBean - private DefaultHousekeeperService housekeeperService; + private HousekeeperService housekeeperService; + @SpyBean + private HousekeeperClient housekeeperClient; @SpyBean private HousekeeperReprocessingService housekeeperReprocessingService; @Autowired @@ -300,7 +303,7 @@ public class HousekeeperServiceTest extends AbstractControllerTest { private void verifyNoRelatedData(EntityId entityId) throws Exception { List expectedTaskTypes = List.of(HousekeeperTaskType.DELETE_TELEMETRY, HousekeeperTaskType.DELETE_ATTRIBUTES, HousekeeperTaskType.DELETE_EVENTS, HousekeeperTaskType.DELETE_ENTITY_ALARMS); for (HousekeeperTaskType taskType : expectedTaskTypes) { - verify(housekeeperService).submitTask(argThat(task -> task.getTaskType() == taskType && task.getEntityId().equals(entityId))); + verify(housekeeperClient).submitTask(argThat(task -> task.getTaskType() == taskType && task.getEntityId().equals(entityId))); } assertThat(getLatestTelemetry(entityId)).isNull(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/HousekeeperService.java b/common/message/src/main/java/org/thingsboard/server/common/msg/housekeeper/HousekeeperClient.java similarity index 88% rename from dao/src/main/java/org/thingsboard/server/dao/housekeeper/HousekeeperService.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/housekeeper/HousekeeperClient.java index e3c0800f69..e2a3af9242 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/HousekeeperService.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/housekeeper/HousekeeperClient.java @@ -13,11 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.housekeeper; +package org.thingsboard.server.common.msg.housekeeper; import org.thingsboard.server.common.data.housekeeper.HousekeeperTask; -public interface HousekeeperService { +public interface HousekeeperClient { void submitTask(HousekeeperTask task); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java b/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java new file mode 100644 index 0000000000..b766f3826b --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java @@ -0,0 +1,58 @@ +/** + * Copyright © 2016-2024 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.housekeeper; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.housekeeper.HousekeeperTask; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.common.msg.housekeeper.HousekeeperClient; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; +import org.thingsboard.server.queue.TbQueueProducer; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.provider.TbQueueProducerProvider; + +@Service +@Slf4j +public class DefaultHousekeeperClient implements HousekeeperClient { + + private final TbQueueProducer> producer; + private final TopicPartitionInfo submitTpi; + + public DefaultHousekeeperClient(TbQueueProducerProvider producerProvider) { + this.producer = producerProvider.getHousekeeperMsgProducer(); + this.submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); + } + + @Override + public void submitTask(HousekeeperTask task) { + log.trace("[{}][{}][{}] Submitting task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType()); + /* + * using msg key as entity id so that msgs related to certain entity are pushed to same partition, + * e.g. on tenant deletion (entity id is tenant id), we need to clean up tenant entities in certain order + * */ + producer.send(submitTpi, new TbProtoQueueMsg<>(task.getEntityId().getId(), ToHousekeeperServiceMsg.newBuilder() + .setTask(TransportProtos.HousekeeperTaskProto.newBuilder() + .setValue(JacksonUtil.toString(task)) + .setTs(task.getTs()) + .setAttempt(0) + .build()) + .build()), null); + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java index 1217186410..e60594c4a1 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java @@ -229,22 +229,24 @@ public class AwsSqsMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEng @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java index 1a4f92fcd3..eb3706657b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java @@ -208,22 +208,24 @@ public class AwsSqsTbCoreQueueFactory implements TbCoreQueueFactory { @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java index 6d48ceb1b6..1653955ae9 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java @@ -33,8 +33,8 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.settings.TbQueueCoreSettings; import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings; import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings; @@ -159,6 +159,10 @@ public class AwsSqsTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory return new TbAwsSqsProducerTemplate<>(otaAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } @PreDestroy private void destroy() { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java index 3aec2e3a39..6801703231 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java @@ -80,6 +80,11 @@ public class AwsSqsTbVersionControlQueueFactory implements TbVersionControlQueue ); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java index c27ad4e092..96a9686d13 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java @@ -131,6 +131,11 @@ public class AwsSqsTransportQueueFactory implements TbTransportQueueFactory { return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/HousekeeperClientQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/HousekeeperClientQueueFactory.java new file mode 100644 index 0000000000..79c38a04e9 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/HousekeeperClientQueueFactory.java @@ -0,0 +1,26 @@ +/** + * Copyright © 2016-2024 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.provider; + +import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; +import org.thingsboard.server.queue.TbQueueProducer; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; + +public interface HousekeeperClientQueueFactory { + + TbQueueProducer> createHousekeeperMsgProducer(); + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java index f08a822466..4ee79aef60 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java @@ -120,4 +120,9 @@ public class InMemoryTbTransportQueueFactory implements TbTransportQueueFactory return new InMemoryTbQueueProducer<>(storage, topicService.buildTopicName(coreSettings.getUsageStatsTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new InMemoryTbQueueProducer<>(storage, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index 63798f95a1..869f364970 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -72,6 +72,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { private final TbQueueAdmin jsExecutorResponseAdmin; private final TbQueueAdmin notificationAdmin; private final TbQueueAdmin fwUpdatesAdmin; + private final TbQueueAdmin housekeeperAdmin; private final AtomicLong consumerCount = new AtomicLong(); public KafkaTbRuleEngineQueueFactory(TopicService topicService, TbKafkaSettings kafkaSettings, @@ -97,6 +98,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { this.jsExecutorResponseAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getJsExecutorResponseConfigs()); this.notificationAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getNotificationsConfigs()); this.fwUpdatesAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getFwUpdatesConfigs()); + this.housekeeperAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getHousekeeperConfigs()); } @Override @@ -230,6 +232,16 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { return requestBuilder.build(); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return TbKafkaProducerTemplate.>builder() + .settings(kafkaSettings) + .clientId("tb-rule-engine-housekeeper-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName(coreSettings.getHousekeeperTopic())) + .admin(housekeeperAdmin) + .build(); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java index 7cecc94290..4140925775 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java @@ -66,6 +66,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { private final TbQueueAdmin transportApiRequestAdmin; private final TbQueueAdmin transportApiResponseAdmin; private final TbQueueAdmin notificationAdmin; + private final TbQueueAdmin housekeeperAdmin; public KafkaTbTransportQueueFactory(TbKafkaSettings kafkaSettings, TbServiceInfoProvider serviceInfoProvider, @@ -90,6 +91,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { this.transportApiRequestAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTransportApiRequestConfigs()); this.transportApiResponseAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTransportApiResponseConfigs()); this.notificationAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getNotificationsConfigs()); + this.housekeeperAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getHousekeeperConfigs()); } @Override @@ -173,6 +175,16 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { return requestBuilder.build(); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return TbKafkaProducerTemplate.>builder() + .settings(kafkaSettings) + .clientId("tb-transport-housekeeper-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName(coreSettings.getHousekeeperTopic())) + .admin(housekeeperAdmin) + .build(); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbVersionControlQueueFactory.java index e39da6bdfc..95086e04b4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbVersionControlQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbVersionControlQueueFactory.java @@ -17,6 +17,7 @@ package org.thingsboard.server.queue.provider; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg; @@ -51,6 +52,7 @@ public class KafkaTbVersionControlQueueFactory implements TbVersionControlQueueF private final TbQueueAdmin coreAdmin; private final TbQueueAdmin vcAdmin; private final TbQueueAdmin notificationAdmin; + private final TbQueueAdmin housekeeperAdmin; public KafkaTbVersionControlQueueFactory(TbKafkaSettings kafkaSettings, TbServiceInfoProvider serviceInfoProvider, @@ -69,6 +71,7 @@ public class KafkaTbVersionControlQueueFactory implements TbVersionControlQueueF this.coreAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCoreConfigs()); this.vcAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getVcConfigs()); this.notificationAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getNotificationsConfigs()); + this.housekeeperAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getHousekeeperConfigs()); } @@ -105,6 +108,16 @@ public class KafkaTbVersionControlQueueFactory implements TbVersionControlQueueF return requestBuilder.build(); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return TbKafkaProducerTemplate.>builder() + .settings(kafkaSettings) + .clientId("tb-vc-housekeeper-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName(coreSettings.getHousekeeperTopic())) + .admin(housekeeperAdmin) + .build(); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java index e6a23b8068..0f2485b6a7 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java @@ -41,8 +41,8 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.pubsub.TbPubSubAdmin; import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate; import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate; @@ -230,22 +230,24 @@ public class PubSubMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEng @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java index 1b52a2ba42..00a6d83bdd 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java @@ -39,8 +39,8 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.pubsub.TbPubSubAdmin; import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate; import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate; @@ -201,22 +201,24 @@ public class PubSubTbCoreQueueFactory implements TbCoreQueueFactory { @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java index 070384e60e..d0192d3094 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java @@ -36,8 +36,8 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.pubsub.TbPubSubAdmin; import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate; import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate; @@ -161,6 +161,11 @@ public class PubSubTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbVersionControlQueueFactory.java index ed40a43deb..509cf371b8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbVersionControlQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbVersionControlQueueFactory.java @@ -79,6 +79,11 @@ public class PubSubTbVersionControlQueueFactory implements TbVersionControlQueue ); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java index 3323308bef..d91dd91873 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueFactory.java @@ -131,6 +131,11 @@ public class PubSubTransportQueueFactory implements TbTransportQueueFactory { return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqMonolithQueueFactory.java index fa89956b64..78b60bca28 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqMonolithQueueFactory.java @@ -41,8 +41,8 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.rabbitmq.TbRabbitMqAdmin; import org.thingsboard.server.queue.rabbitmq.TbRabbitMqConsumerTemplate; import org.thingsboard.server.queue.rabbitmq.TbRabbitMqProducerTemplate; @@ -227,22 +227,24 @@ public class RabbitMqMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbRabbitMqConsumerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbRabbitMqConsumerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbCoreQueueFactory.java index b1f3de672f..db3414b62d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbCoreQueueFactory.java @@ -38,8 +38,8 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.rabbitmq.TbRabbitMqAdmin; import org.thingsboard.server.queue.rabbitmq.TbRabbitMqConsumerTemplate; import org.thingsboard.server.queue.rabbitmq.TbRabbitMqProducerTemplate; @@ -200,22 +200,24 @@ public class RabbitMqTbCoreQueueFactory implements TbCoreQueueFactory { @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbRabbitMqConsumerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbRabbitMqConsumerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbRuleEngineQueueFactory.java index 929ea8298c..556aa968c6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbRuleEngineQueueFactory.java @@ -160,6 +160,11 @@ public class RabbitMqTbRuleEngineQueueFactory implements TbRuleEngineQueueFactor return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbVersionControlQueueFactory.java index bd14db38b4..a614df3fdb 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbVersionControlQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTbVersionControlQueueFactory.java @@ -79,6 +79,11 @@ public class RabbitMqTbVersionControlQueueFactory implements TbVersionControlQue ); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java index d5fb7bd0e2..389a6d11ee 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/RabbitMqTransportQueueFactory.java @@ -132,6 +132,11 @@ public class RabbitMqTransportQueueFactory implements TbTransportQueueFactory { return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbRabbitMqProducerTemplate<>(coreAdmin, rabbitMqSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusMonolithQueueFactory.java index 5f4be98ca0..a1f1c7adc8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusMonolithQueueFactory.java @@ -45,8 +45,8 @@ import org.thingsboard.server.queue.azure.servicebus.TbServiceBusSettings; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.settings.TbQueueCoreSettings; import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings; import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings; @@ -226,22 +226,24 @@ public class ServiceBusMonolithQueueFactory implements TbCoreQueueFactory, TbRul @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbServiceBusConsumerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbServiceBusConsumerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbCoreQueueFactory.java index 579fba272c..854295cb6b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbCoreQueueFactory.java @@ -44,8 +44,8 @@ import org.thingsboard.server.queue.azure.servicebus.TbServiceBusSettings; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.settings.TbQueueCoreSettings; import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings; import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings; @@ -201,22 +201,24 @@ public class ServiceBusTbCoreQueueFactory implements TbCoreQueueFactory { @Override public TbQueueProducer> createHousekeeperMsgProducer() { - return null; + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); } @Override public TbQueueConsumer> createHousekeeperMsgConsumer() { - return null; + return new TbServiceBusConsumerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @Override public TbQueueProducer> createHousekeeperReprocessingMsgProducer() { - return null; + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic())); } @Override public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() { - return null; + return new TbServiceBusConsumerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()), + msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders())); } @PreDestroy diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbRuleEngineQueueFactory.java index b09ab1d30d..066e84e780 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbRuleEngineQueueFactory.java @@ -41,8 +41,8 @@ import org.thingsboard.server.queue.azure.servicebus.TbServiceBusSettings; import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.settings.TbQueueCoreSettings; import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings; import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings; @@ -160,6 +160,11 @@ public class ServiceBusTbRuleEngineQueueFactory implements TbRuleEngineQueueFact return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbVersionControlQueueFactory.java index 2e5e5c4bc3..bed6c427c9 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbVersionControlQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTbVersionControlQueueFactory.java @@ -79,6 +79,11 @@ public class ServiceBusTbVersionControlQueueFactory implements TbVersionControlQ ); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java index d782a1ff9e..0db62ac90c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/ServiceBusTransportQueueFactory.java @@ -133,6 +133,11 @@ public class ServiceBusTransportQueueFactory implements TbTransportQueueFactory return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic())); } + @Override + public TbQueueProducer> createHousekeeperMsgProducer() { + return new TbServiceBusProducerTemplate<>(coreAdmin, serviceBusSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic())); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java index 866eb931f2..20a67c2e09 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java @@ -37,7 +37,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; * Responsible for initialization of various Producers and Consumers used by TB Core Node. * Implementation Depends on the queue queue.type from yml or TB_QUEUE_TYPE environment variable */ -public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory { +public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory { /** * Used to push messages to instances of TB Transport Service @@ -132,8 +132,6 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory { */ TbQueueProducer> createVersionControlMsgProducer(); - TbQueueProducer> createHousekeeperMsgProducer(); - TbQueueConsumer> createHousekeeperMsgConsumer(); TbQueueProducer> createHousekeeperReprocessingMsgProducer(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java index e6bc8f0484..aa2c55fca1 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java @@ -43,7 +43,6 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider { private TbQueueProducer> toUsageStats; private TbQueueProducer> toVersionControl; private TbQueueProducer> toHousekeeper; - private TbQueueProducer> toHousekeeperReprocessing; public TbCoreQueueProducerProvider(TbCoreQueueFactory tbQueueProvider) { this.tbQueueProvider = tbQueueProvider; @@ -59,7 +58,6 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider { this.toUsageStats = tbQueueProvider.createToUsageStatsServiceMsgProducer(); this.toVersionControl = tbQueueProvider.createVersionControlMsgProducer(); this.toHousekeeper = tbQueueProvider.createHousekeeperMsgProducer(); - this.toHousekeeperReprocessing = tbQueueProvider.createHousekeeperReprocessingMsgProducer(); } @Override @@ -102,9 +100,4 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider { return toHousekeeper; } - @Override - public TbQueueProducer> getHousekeeperReprocessingMsgProducer() { - return toHousekeeperReprocessing; - } - } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java index ff56d2e8be..a752301cd8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java @@ -82,6 +82,4 @@ public interface TbQueueProducerProvider { TbQueueProducer> getHousekeeperMsgProducer(); - TbQueueProducer> getHousekeeperReprocessingMsgProducer(); - } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java index ee6700906c..748c48e52a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java @@ -20,6 +20,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg; +import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; @@ -40,6 +41,7 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider { private TbQueueProducer> toRuleEngineNotifications; private TbQueueProducer> toTbCoreNotifications; private TbQueueProducer> toUsageStats; + private TbQueueProducer> toHousekeeper; public TbRuleEngineProducerProvider(TbRuleEngineQueueFactory tbQueueProvider) { this.tbQueueProvider = tbQueueProvider; @@ -53,6 +55,7 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider { this.toRuleEngineNotifications = tbQueueProvider.createRuleEngineNotificationsMsgProducer(); this.toTbCoreNotifications = tbQueueProvider.createTbCoreNotificationsMsgProducer(); this.toUsageStats = tbQueueProvider.createToUsageStatsServiceMsgProducer(); + this.toHousekeeper = tbQueueProvider.createHousekeeperMsgProducer(); } @Override @@ -91,13 +94,8 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider { } @Override - public TbQueueProducer> getHousekeeperMsgProducer() { - return null; // fixme - } - - @Override - public TbQueueProducer> getHousekeeperReprocessingMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Rule Engine!"); + public TbQueueProducer> getHousekeeperMsgProducer() { + return toHousekeeper; } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java index cd4575af6a..f23c7e47f3 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java @@ -32,7 +32,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; * Responsible for initialization of various Producers and Consumers used by TB Core Node. * Implementation Depends on the queue queue.type from yml or TB_QUEUE_TYPE environment variable */ -public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory { +public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory { /** * Used to push messages to instances of TB Transport Service diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java index 909089f175..c1a0a0fd24 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueFactory.java @@ -26,7 +26,7 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -public interface TbTransportQueueFactory extends TbUsageStatsClientQueueFactory { +public interface TbTransportQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory { TbQueueRequestTemplate, TbProtoQueueMsg> createTransportApiRequestTemplate(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java index 869c588fc7..c0472ba6f0 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java @@ -20,6 +20,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg; +import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; @@ -38,6 +39,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider private TbQueueProducer> toTbCore; private TbQueueProducer> toTbCoreNotifications; private TbQueueProducer> toUsageStats; + private TbQueueProducer> toHousekeeper; public TbTransportQueueProducerProvider(TbTransportQueueFactory tbQueueProvider) { this.tbQueueProvider = tbQueueProvider; @@ -49,6 +51,7 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider this.toRuleEngine = tbQueueProvider.createRuleEngineMsgProducer(); this.toUsageStats = tbQueueProvider.createToUsageStatsServiceMsgProducer(); this.toTbCoreNotifications = tbQueueProvider.createTbCoreNotificationsMsgProducer(); + this.toHousekeeper = tbQueueProvider.createHousekeeperMsgProducer(); } @Override @@ -87,12 +90,8 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider } @Override - public TbQueueProducer> getHousekeeperMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Transport!"); + public TbQueueProducer> getHousekeeperMsgProducer() { + return toHousekeeper; } - @Override - public TbQueueProducer> getHousekeeperReprocessingMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Transport!"); - } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java index 3f4bd209e4..c4e8b697a3 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java @@ -20,6 +20,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg; +import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; @@ -36,6 +37,7 @@ public class TbVersionControlProducerProvider implements TbQueueProducerProvider private final TbVersionControlQueueFactory tbQueueProvider; private TbQueueProducer> toTbCoreNotifications; private TbQueueProducer> toUsageStats; + private TbQueueProducer> toHousekeeper; public TbVersionControlProducerProvider(TbVersionControlQueueFactory tbQueueProvider) { this.tbQueueProvider = tbQueueProvider; @@ -45,6 +47,7 @@ public class TbVersionControlProducerProvider implements TbQueueProducerProvider public void init() { this.toTbCoreNotifications = tbQueueProvider.createTbCoreNotificationsMsgProducer(); this.toUsageStats = tbQueueProvider.createToUsageStatsServiceMsgProducer(); + this.toHousekeeper = tbQueueProvider.createHousekeeperMsgProducer(); } @Override @@ -54,7 +57,7 @@ public class TbVersionControlProducerProvider implements TbQueueProducerProvider @Override public TbQueueProducer> getRuleEngineMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); + throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); } @Override @@ -83,12 +86,8 @@ public class TbVersionControlProducerProvider implements TbQueueProducerProvider } @Override - public TbQueueProducer> getHousekeeperMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); + public TbQueueProducer> getHousekeeperMsgProducer() { + return toHousekeeper; } - @Override - public TbQueueProducer> getHousekeeperReprocessingMsgProducer() { - throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); - } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlQueueFactory.java index 168a466015..3917df24cd 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlQueueFactory.java @@ -25,7 +25,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; * Responsible for initialization of various Producers and Consumers used by TB Version Control Node. * Implementation Depends on the queue queue.type from yml or TB_QUEUE_TYPE environment variable */ -public interface TbVersionControlQueueFactory extends TbUsageStatsClientQueueFactory { +public interface TbVersionControlQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory { /** * Used to push notifications to other instances of TB Core Service diff --git a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java index 015917e592..56946c4705 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java @@ -17,8 +17,6 @@ package org.thingsboard.server.dao.housekeeper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionalEventListener; import org.thingsboard.server.common.data.EntityType; @@ -26,23 +24,17 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.housekeeper.HousekeeperTask; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.dao.entity.EntityDaoService; -import org.thingsboard.server.dao.entity.EntityServiceRegistry; +import org.thingsboard.server.common.msg.housekeeper.HousekeeperClient; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.relation.RelationService; -import java.util.Optional; - @Component @RequiredArgsConstructor @Slf4j public class CleanUpService { - private final Optional housekeeperService; + private final HousekeeperClient housekeeperClient; private final RelationService relationService; - @Autowired - @Lazy - private EntityServiceRegistry entityServiceRegistry; @TransactionalEventListener(fallbackExecution = true) public void handleEntityDeletionEvent(DeleteEntityEvent event) { @@ -51,9 +43,7 @@ public class CleanUpService { log.trace("[{}][{}][{}] Handling entity deletion event", tenantId, entityId.getEntityType(), entityId.getId()); cleanUpRelatedData(tenantId, entityId); if (entityId.getEntityType() == EntityType.USER) { - housekeeperService.ifPresent(housekeeperService -> { - housekeeperService.submitTask(HousekeeperTask.unassignAlarms((User) event.getEntity())); - }); + housekeeperClient.submitTask(HousekeeperTask.unassignAlarms((User) event.getEntity())); } } @@ -61,25 +51,16 @@ public class CleanUpService { log.debug("[{}][{}][{}] Cleaning up related data", tenantId, entityId.getEntityType(), entityId.getId()); // todo: skipped entities list relationService.deleteEntityRelations(tenantId, entityId); - housekeeperService.ifPresent(housekeeperService -> { - housekeeperService.submitTask(HousekeeperTask.deleteAttributes(tenantId, entityId)); - housekeeperService.submitTask(HousekeeperTask.deleteTelemetry(tenantId, entityId)); - housekeeperService.submitTask(HousekeeperTask.deleteEvents(tenantId, entityId)); - housekeeperService.submitTask(HousekeeperTask.deleteEntityAlarms(tenantId, entityId)); - }); + housekeeperClient.submitTask(HousekeeperTask.deleteAttributes(tenantId, entityId)); + housekeeperClient.submitTask(HousekeeperTask.deleteTelemetry(tenantId, entityId)); + housekeeperClient.submitTask(HousekeeperTask.deleteEvents(tenantId, entityId)); + housekeeperClient.submitTask(HousekeeperTask.deleteEntityAlarms(tenantId, entityId)); } public void removeTenantEntities(TenantId tenantId, EntityType... entityTypes) { - housekeeperService.ifPresentOrElse(housekeeperService -> { - for (EntityType entityType : entityTypes) { - housekeeperService.submitTask(HousekeeperTask.deleteEntities(tenantId, entityType)); - } - }, () -> { - for (EntityType entityType : entityTypes) { - EntityDaoService entityService = entityServiceRegistry.getServiceByEntityType(entityType); - entityService.deleteByTenantId(tenantId); - } - }); + for (EntityType entityType : entityTypes) { + housekeeperClient.submitTask(HousekeeperTask.deleteEntities(tenantId, entityType)); + } } } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java index e8232af692..f222c05b60 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java @@ -42,6 +42,8 @@ import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileTra import org.thingsboard.server.common.data.device.profile.DeviceProfileData; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; +import org.thingsboard.server.common.data.housekeeper.EntitiesDeletionHousekeeperTask; +import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; @@ -51,6 +53,9 @@ import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.dao.audit.AuditLogLevelFilter; import org.thingsboard.server.dao.audit.AuditLogLevelMask; import org.thingsboard.server.dao.audit.AuditLogLevelProperties; +import org.thingsboard.server.dao.entity.EntityDaoService; +import org.thingsboard.server.dao.entity.EntityServiceRegistry; +import org.thingsboard.server.common.msg.housekeeper.HousekeeperClient; import org.thingsboard.server.dao.tenant.TenantService; import java.io.IOException; @@ -77,6 +82,9 @@ public abstract class AbstractServiceTest { @Autowired protected TenantService tenantService; + @Autowired + protected EntityServiceRegistry entityServiceRegistry; + protected TenantId tenantId; @Before @@ -126,6 +134,16 @@ public abstract class AbstractServiceTest { return new AuditLogLevelFilter(props); } + @Bean + public HousekeeperClient housekeeperClient() { + return task -> { + if (task.getTaskType() == HousekeeperTaskType.DELETE_ENTITIES) { + EntityDaoService entityService = entityServiceRegistry.getServiceByEntityType(((EntitiesDeletionHousekeeperTask) task).getEntityType()); + entityService.deleteByTenantId(task.getTenantId()); + } + }; + } + protected DeviceProfile createDeviceProfile(TenantId tenantId, String name) { DeviceProfile deviceProfile = new DeviceProfile(); deviceProfile.setTenantId(tenantId);