From 25f82d384b8d4d1a1847016360388688a1750d48 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Wed, 16 May 2018 20:58:46 +0300 Subject: [PATCH] Message queue limit per tenant. Message queue periodic cleanup. --- .../server/actors/ActorSystemContext.java | 3 +- .../RuleChainActorMessageProcessor.java | 8 +- .../actors/shared/ComponentMsgProcessor.java | 7 +- .../service/queue/DefaultMsgQueueService.java | 102 +++++++++++++++ .../server/service/queue/MsgQueueService.java | 32 +++++ .../src/main/resources/thingsboard.yml | 13 +- .../AbstractRuleEngineControllerTest.java | 3 +- ...AbstractRuleEngineFlowIntegrationTest.java | 6 +- ...actRuleEngineLifecycleIntegrationTest.java | 2 +- .../server/dao/queue/MsgQueue.java | 10 +- .../server/dao/queue/QueueBenchmark.java | 4 +- .../db/nosql}/CassandraMsgQueue.java | 21 +++- .../cassandra => queue/db/nosql}/MsgAck.java | 2 +- .../db/nosql}/QueuePartitioner.java | 4 +- .../db/nosql}/UnprocessedMsgFilter.java | 2 +- .../repository}/CassandraAckRepository.java | 6 +- .../repository}/CassandraMsgRepository.java | 4 +- ...CassandraProcessedPartitionRepository.java | 4 +- .../db}/repository/AckRepository.java | 4 +- .../db}/repository/MsgRepository.java | 2 +- .../ProcessedPartitionRepository.java | 2 +- .../queue/{jpa => db/sql}/SqlMsgQueue.java | 2 +- .../memory}/InMemoryMsgKey.java | 2 +- .../dao/queue/memory/InMemoryMsgQueue.java | 118 ++++++++++++++++++ .../dao/sql/queue/InMemoryMsgQueue.java | 114 ----------------- .../db/nosql}/QueuePartitionerTest.java | 5 +- .../db/nosql}/UnprocessedMsgFilterTest.java | 4 +- .../CassandraAckRepositoryTest.java | 5 +- .../CassandraMsgRepositoryTest.java | 3 +- ...andraProcessedPartitionRepositoryTest.java | 3 +- 30 files changed, 328 insertions(+), 169 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java create mode 100644 application/src/main/java/org/thingsboard/server/service/queue/MsgQueueService.java rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db/nosql}/CassandraMsgQueue.java (73%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db/nosql}/MsgAck.java (93%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db/nosql}/QueuePartitioner.java (95%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db/nosql}/UnprocessedMsgFilter.java (95%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra/repository/impl => queue/db/nosql/repository}/CassandraAckRepository.java (91%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra/repository/impl => queue/db/nosql/repository}/CassandraMsgRepository.java (94%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra/repository/impl => queue/db/nosql/repository}/CassandraProcessedPartitionRepository.java (93%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db}/repository/AckRepository.java (86%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db}/repository/MsgRepository.java (93%) rename dao/src/main/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db}/repository/ProcessedPartitionRepository.java (93%) rename dao/src/main/java/org/thingsboard/server/dao/queue/{jpa => db/sql}/SqlMsgQueue.java (93%) rename dao/src/main/java/org/thingsboard/server/dao/{sql/queue => queue/memory}/InMemoryMsgKey.java (94%) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgQueue.java rename dao/src/test/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db/nosql}/QueuePartitionerTest.java (92%) rename dao/src/test/java/org/thingsboard/server/dao/{service/queue/cassandra => queue/db/nosql}/UnprocessedMsgFilterTest.java (89%) rename dao/src/test/java/org/thingsboard/server/dao/{service/queue/cassandra/repository/impl => queue/db/nosql/repository}/CassandraAckRepositoryTest.java (94%) rename dao/src/test/java/org/thingsboard/server/dao/{service/queue/cassandra/repository/impl => queue/db/nosql/repository}/CassandraMsgRepositoryTest.java (97%) rename dao/src/test/java/org/thingsboard/server/dao/{service/queue/cassandra/repository/impl => queue/db/nosql/repository}/CassandraProcessedPartitionRepositoryTest.java (97%) 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 665d2d9489..864059892b 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -63,6 +63,7 @@ import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.ExternalCallExecutorService; import org.thingsboard.server.service.mail.MailExecutorService; +import org.thingsboard.server.service.queue.MsgQueueService; import org.thingsboard.server.service.rpc.DeviceRpcService; import org.thingsboard.server.service.script.JsExecutorService; import org.thingsboard.server.service.state.DeviceStateService; @@ -196,7 +197,7 @@ public class ActorSystemContext { @Autowired @Getter - private MsgQueue msgQueue; + private MsgQueueService msgQueueService; @Autowired @Getter diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index d069cb0eae..dda12e5a52 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -95,12 +95,12 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor ruleNodeList) { for (RuleNode ruleNode : ruleNodeList) { - for (TbMsg tbMsg : queue.findUnprocessed(ruleNode.getId().getId(), 0L)) { + for (TbMsg tbMsg : queue.findUnprocessed(tenantId, ruleNode.getId().getId(), 0L)) { pushMsgToNode(nodeActors.get(ruleNode.getId()), tbMsg, ""); } } if (firstNode != null) { - for (TbMsg tbMsg : queue.findUnprocessed(entityId.getId(), 0L)) { + for (TbMsg tbMsg : queue.findUnprocessed(tenantId, entityId.getId(), 0L)) { pushMsgToNode(firstNode, tbMsg, ""); } } @@ -215,7 +215,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor extends Abstract protected final TenantId tenantId; protected final T entityId; - protected final MsgQueue queue; + protected final MsgQueueService queue; protected ComponentLifecycleState state; protected ComponentMsgProcessor(ActorSystemContext systemContext, LoggingAdapter logger, TenantId tenantId, T id) { super(systemContext, logger); this.tenantId = tenantId; this.entityId = id; - this.queue = systemContext.getMsgQueue(); + this.queue = systemContext.getMsgQueueService(); } public abstract void start(ActorContext context) throws Exception; @@ -88,7 +89,7 @@ public abstract class ComponentMsgProcessor extends Abstract protected void putToQueue(final TbMsg tbMsg, final Consumer onSuccess) { EntityId entityId = tbMsg.getRuleNodeId() != null ? tbMsg.getRuleNodeId() : tbMsg.getRuleChainId(); - Futures.addCallback(queue.put(tbMsg, entityId.getId(), 0), new FutureCallback() { + Futures.addCallback(queue.put(this.tenantId, tbMsg, entityId.getId(), 0), new FutureCallback() { @Override public void onSuccess(@Nullable Void result) { onSuccess.accept(tbMsg); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java new file mode 100644 index 0000000000..a67278cc6d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java @@ -0,0 +1,102 @@ +/** + * Copyright © 2016-2018 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; + +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.dao.queue.MsgQueue; + +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +@Service +@Slf4j +public class DefaultMsgQueueService implements MsgQueueService { + + @Value("${rule.queue.max_size}") + private long queueMaxSize; + + @Value("${rule.queue.cleanup_period}") + private long queueCleanUpPeriod; + + @Autowired + private MsgQueue msgQueue; + + private ScheduledExecutorService cleanupExecutor; + + private Map pendingCountPerTenant = new ConcurrentHashMap<>(); + + @PostConstruct + public void init() { + if (queueCleanUpPeriod > 0) { + cleanupExecutor = Executors.newSingleThreadScheduledExecutor(); + cleanupExecutor.scheduleAtFixedRate(() -> cleanup(), + queueCleanUpPeriod, queueCleanUpPeriod, TimeUnit.SECONDS); + } + } + + @PreDestroy + public void stop() { + if (cleanupExecutor != null) { + cleanupExecutor.shutdownNow(); + } + } + + @Override + public ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { + AtomicLong pendingMsgCount = pendingCountPerTenant.computeIfAbsent(tenantId, key -> new AtomicLong()); + if (pendingMsgCount.incrementAndGet() < queueMaxSize) { + return msgQueue.put(tenantId, msg, nodeId, clusterPartition); + } else { + pendingMsgCount.decrementAndGet(); + return Futures.immediateFailedFuture(new RuntimeException("Message queue is full!")); + } + } + + @Override + public ListenableFuture ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { + ListenableFuture result = msgQueue.ack(tenantId, msg, nodeId, clusterPartition); + AtomicLong pendingMsgCount = pendingCountPerTenant.computeIfAbsent(tenantId, key -> new AtomicLong()); + pendingMsgCount.decrementAndGet(); + return result; + } + + @Override + public Iterable findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { + return msgQueue.findUnprocessed(tenantId, nodeId, clusterPartition); + } + + private void cleanup() { + pendingCountPerTenant.forEach((tenantId, pendingMsgCount) -> { + pendingMsgCount.set(0); + msgQueue.cleanUp(tenantId); + }); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/queue/MsgQueueService.java b/application/src/main/java/org/thingsboard/server/service/queue/MsgQueueService.java new file mode 100644 index 0000000000..2cf001ba73 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/MsgQueueService.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2018 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; + +import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.TbMsg; + +import java.util.UUID; + +public interface MsgQueueService { + + ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition); + + ListenableFuture ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition); + + Iterable findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition); + +} diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index ac7a25923d..dda6e1595f 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -305,17 +305,20 @@ spring: rule: queue: - type: "memory" - max_size: 10000 - + #Message queue type (memory or db) + type: "${RULE_QUEUE_TYPE:memory}" + #Message queue maximum size (per tenant) + max_size: "${RULE_QUEUE_MAX_SIZE:100}" + #Message queue cleanup period in seconds + cleanup_period: "${RULE_QUEUE_CLEANUP_PERIOD:3600}" # PostgreSQL DAO Configuration #spring: # data: -# jpa: +# sql: # repositories: # enabled: "true" -# jpa: +# sql: # hibernate: # ddl-auto: "validate" # database-platform: "${SPRING_JPA_DATABASE_PLATFORM:org.hibernate.dialect.PostgreSQLDialect}" diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java index 5b895b1513..6c2d3dbf7d 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.dao.queue.MsgQueue; import org.thingsboard.server.dao.rule.RuleChainService; +import org.thingsboard.server.service.queue.MsgQueueService; import java.io.IOException; @@ -41,7 +42,7 @@ public class AbstractRuleEngineControllerTest extends AbstractControllerTest { protected RuleChainService ruleChainService; @Autowired - protected MsgQueue msgQueue; + protected MsgQueueService msgQueueService; protected RuleChain saveRuleChain(RuleChain ruleChain) throws Exception { return doPost("/api/ruleChain", ruleChain, RuleChain.class); diff --git a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java index a294816f49..356dfeea41 100644 --- a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java @@ -189,7 +189,7 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule Assert.assertEquals("serverAttributeValue1", getMetadata(outEvent).get("ss_serverAttributeKey1").asText()); Assert.assertEquals("serverAttributeValue2", getMetadata(outEvent).get("ss_serverAttributeKey2").asText()); - List unAckMsgList = Lists.newArrayList(msgQueue.findUnprocessed(ruleChain.getId().getId(), 0L)); + List unAckMsgList = Lists.newArrayList(msgQueueService.findUnprocessed(savedTenant.getId(), ruleChain.getId().getId(), 0L)); Assert.assertEquals(0, unAckMsgList.size()); } @@ -306,10 +306,10 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule Assert.assertEquals("serverAttributeValue1", getMetadata(outEvent).get("ss_serverAttributeKey1").asText()); Assert.assertEquals("serverAttributeValue2", getMetadata(outEvent).get("ss_serverAttributeKey2").asText()); - List unAckMsgList = Lists.newArrayList(msgQueue.findUnprocessed(rootRuleChain.getId().getId(), 0L)); + List unAckMsgList = Lists.newArrayList(msgQueueService.findUnprocessed(savedTenant.getId(), rootRuleChain.getId().getId(), 0L)); Assert.assertEquals(0, unAckMsgList.size()); - unAckMsgList = Lists.newArrayList(msgQueue.findUnprocessed(secondaryRuleChain.getId().getId(), 0L)); + unAckMsgList = Lists.newArrayList(msgQueueService.findUnprocessed(savedTenant.getId(), secondaryRuleChain.getId().getId(), 0L)); Assert.assertEquals(0, unAckMsgList.size()); } diff --git a/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java index 0ea6ff4aec..2f25b97eee 100644 --- a/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java @@ -186,7 +186,7 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac new TbMsgMetaData(), "{}", ruleChain.getId(), null, 0L); - msgQueue.put(tbMsg, ruleChain.getId().getId(), 0L); + msgQueueService.put(device.getTenantId(), tbMsg, ruleChain.getId().getId(), 0L); Thread.sleep(1000); diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/MsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/MsgQueue.java index e49b89eb81..19021eb95c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/MsgQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/MsgQueue.java @@ -16,15 +16,19 @@ package org.thingsboard.server.dao.queue; import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; import java.util.UUID; public interface MsgQueue { - ListenableFuture put(TbMsg msg, UUID nodeId, long clusterPartition); + ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition); - ListenableFuture ack(TbMsg msg, UUID nodeId, long clusterPartition); + ListenableFuture ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition); + + Iterable findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition); + + ListenableFuture cleanUp(TenantId tenantId); - Iterable findUnprocessed(UUID nodeId, long clusterPartition); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/QueueBenchmark.java b/dao/src/main/java/org/thingsboard/server/dao/queue/QueueBenchmark.java index ef55bcb1de..ca61a63a4c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/QueueBenchmark.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/QueueBenchmark.java @@ -28,8 +28,10 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.CommandLineRunner; import org.springframework.boot.SpringApplication; import org.springframework.context.annotation.Bean; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; @@ -80,7 +82,7 @@ public class QueueBenchmark implements CommandLineRunner { try { TbMsg msg = randomMsg(); UUID nodeId = UUIDs.timeBased(); - ListenableFuture put = msgQueue.put(msg, nodeId, 100L); + ListenableFuture put = msgQueue.put(new TenantId(EntityId.NULL_UUID), msg, nodeId, 100L); // ListenableFuture put = msgQueue.ack(msg, nodeId, 100L); Futures.addCallback(put, new FutureCallback() { @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/CassandraMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java similarity index 73% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/CassandraMsgQueue.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java index 361b18ce4b..ceaab586cf 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/CassandraMsgQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java @@ -13,24 +13,28 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra; +package org.thingsboard.server.dao.queue.db.nosql; import com.datastax.driver.core.utils.UUIDs; import com.google.common.collect.Lists; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.dao.queue.MsgQueue; -import org.thingsboard.server.dao.service.queue.cassandra.repository.AckRepository; -import org.thingsboard.server.dao.service.queue.cassandra.repository.MsgRepository; +import org.thingsboard.server.dao.queue.db.repository.AckRepository; +import org.thingsboard.server.dao.queue.db.repository.MsgRepository; import org.thingsboard.server.dao.util.NoSqlDao; import java.util.List; import java.util.UUID; @Component +@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "db") @Slf4j @NoSqlDao public class CassandraMsgQueue implements MsgQueue { @@ -45,21 +49,21 @@ public class CassandraMsgQueue implements MsgQueue { private QueuePartitioner queuePartitioner; @Override - public ListenableFuture put(TbMsg msg, UUID nodeId, long clusterPartition) { + public ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { long msgTime = getMsgTime(msg); long tsPartition = queuePartitioner.getPartition(msgTime); return msgRepository.save(msg, nodeId, clusterPartition, tsPartition, msgTime); } @Override - public ListenableFuture ack(TbMsg msg, UUID nodeId, long clusterPartition) { + public ListenableFuture ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { long tsPartition = queuePartitioner.getPartition(getMsgTime(msg)); MsgAck ack = new MsgAck(msg.getId(), nodeId, clusterPartition, tsPartition); return ackRepository.ack(ack); } @Override - public Iterable findUnprocessed(UUID nodeId, long clusterPartition) { + public Iterable findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { List unprocessedMsgs = Lists.newArrayList(); for (Long tsPartition : queuePartitioner.findUnprocessedPartitions(nodeId, clusterPartition)) { List msgs = msgRepository.findMsgs(nodeId, clusterPartition, tsPartition); @@ -69,6 +73,11 @@ public class CassandraMsgQueue implements MsgQueue { return unprocessedMsgs; } + @Override + public ListenableFuture cleanUp(TenantId tenantId) { + return Futures.immediateFuture(null); + } + private long getMsgTime(TbMsg msg) { return UUIDs.unixTimestamp(msg.getId()); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/MsgAck.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/MsgAck.java similarity index 93% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/MsgAck.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/MsgAck.java index fed6e856eb..1b1cd3f467 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/MsgAck.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/MsgAck.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra; +package org.thingsboard.server.dao.queue.db.nosql; import lombok.Data; import lombok.EqualsAndHashCode; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/QueuePartitioner.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java similarity index 95% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/QueuePartitioner.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java index a60f685f29..6076d93e9f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/QueuePartitioner.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java @@ -13,13 +13,13 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra; +package org.thingsboard.server.dao.queue.db.nosql; import com.google.common.collect.Lists; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; -import org.thingsboard.server.dao.service.queue.cassandra.repository.ProcessedPartitionRepository; +import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; import org.thingsboard.server.dao.timeseries.TsPartitionDate; import org.thingsboard.server.dao.util.NoSqlDao; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/UnprocessedMsgFilter.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilter.java similarity index 95% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/UnprocessedMsgFilter.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilter.java index 4dcd351763..c912e8e114 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/UnprocessedMsgFilter.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilter.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra; +package org.thingsboard.server.dao.queue.db.nosql; import org.springframework.stereotype.Component; import org.thingsboard.server.common.msg.TbMsg; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraAckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java similarity index 91% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraAckRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java index 1f62f2bf4e..1ffbec3d81 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraAckRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; +package org.thingsboard.server.dao.queue.db.nosql.repository; import com.datastax.driver.core.*; import com.google.common.base.Function; @@ -22,8 +22,8 @@ import com.google.common.util.concurrent.ListenableFuture; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.service.queue.cassandra.MsgAck; -import org.thingsboard.server.dao.service.queue.cassandra.repository.AckRepository; +import org.thingsboard.server.dao.queue.db.nosql.MsgAck; +import org.thingsboard.server.dao.queue.db.repository.AckRepository; import org.thingsboard.server.dao.util.NoSqlDao; import java.util.ArrayList; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraMsgRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java similarity index 94% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraMsgRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java index 2a70a8909a..2699a4d073 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraMsgRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; +package org.thingsboard.server.dao.queue.db.nosql.repository; import com.datastax.driver.core.*; import com.google.common.base.Function; @@ -23,7 +23,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.service.queue.cassandra.repository.MsgRepository; +import org.thingsboard.server.dao.queue.db.repository.MsgRepository; import org.thingsboard.server.dao.util.NoSqlDao; import java.util.ArrayList; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraProcessedPartitionRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java similarity index 93% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraProcessedPartitionRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java index b0eacfa1b3..60d5b0f5ea 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraProcessedPartitionRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; +package org.thingsboard.server.dao.queue.db.nosql.repository; import com.datastax.driver.core.*; import com.google.common.base.Function; @@ -22,7 +22,7 @@ import com.google.common.util.concurrent.ListenableFuture; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.service.queue.cassandra.repository.ProcessedPartitionRepository; +import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; import org.thingsboard.server.dao.util.NoSqlDao; import java.util.Optional; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/AckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java similarity index 86% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/AckRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java index d7cdb0cc06..458dba81cb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/AckRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java @@ -13,10 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository; +package org.thingsboard.server.dao.queue.db.repository; import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.dao.service.queue.cassandra.MsgAck; +import org.thingsboard.server.dao.queue.db.nosql.MsgAck; import java.util.List; import java.util.UUID; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/MsgRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java similarity index 93% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/MsgRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java index d54f1af59f..0ca6900fc9 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/MsgRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository; +package org.thingsboard.server.dao.queue.db.repository; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.msg.TbMsg; diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/ProcessedPartitionRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java similarity index 93% rename from dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/ProcessedPartitionRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java index a50ab6143e..b11fc6cc1c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/queue/cassandra/repository/ProcessedPartitionRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository; +package org.thingsboard.server.dao.queue.db.repository; import com.google.common.util.concurrent.ListenableFuture; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/jpa/SqlMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java similarity index 93% rename from dao/src/main/java/org/thingsboard/server/dao/queue/jpa/SqlMsgQueue.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java index eb3c9b95f4..f9dda43a36 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/jpa/SqlMsgQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.queue.jpa; +package org.thingsboard.server.dao.queue.db.sql; //@todo-vp: implement public class SqlMsgQueue { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgKey.java b/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgKey.java similarity index 94% rename from dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgKey.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgKey.java index 2090edf381..bb44a2b192 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgKey.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgKey.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.sql.queue; +package org.thingsboard.server.dao.queue.memory; import lombok.Data; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java new file mode 100644 index 0000000000..eecb782a1e --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java @@ -0,0 +1,118 @@ +/** + * Copyright © 2016-2018 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.dao.queue.memory; + +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningExecutorService; +import com.google.common.util.concurrent.MoreExecutors; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.dao.queue.MsgQueue; + +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.*; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Executors; + +/** + * Created by ashvayka on 27.04.18. + */ +@Component +@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "memory", matchIfMissing = true) +@Slf4j +public class InMemoryMsgQueue implements MsgQueue { + + private ListeningExecutorService queueExecutor; + private Map>> data = new HashMap<>(); + + @PostConstruct + public void init() { + // Should be always single threaded due to absence of locks. + queueExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); + } + + @PreDestroy + public void stop() { + if (queueExecutor != null) { + queueExecutor.shutdownNow(); + } + } + + @Override + public ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { + return queueExecutor.submit(() -> { + data.computeIfAbsent(tenantId, key -> new HashMap<>()). + computeIfAbsent(new InMemoryMsgKey(nodeId, clusterPartition), key -> new HashMap<>()).put(msg.getId(), msg); + return null; + }); + } + + @Override + public ListenableFuture ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { + return queueExecutor.submit(() -> { + Map> tenantMap = data.get(tenantId); + if (tenantMap != null) { + InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); + Map map = tenantMap.get(key); + if (map != null) { + map.remove(msg.getId()); + if (map.isEmpty()) { + tenantMap.remove(key); + } + } + if (tenantMap.isEmpty()) { + data.remove(tenantId); + } + } + return null; + }); + } + + @Override + public Iterable findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { + ListenableFuture> list = queueExecutor.submit(() -> { + Map> tenantMap = data.get(tenantId); + if (tenantMap != null) { + InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); + Map map = tenantMap.get(key); + if (map != null) { + return new ArrayList<>(map.values()); + } else { + return Collections.emptyList(); + } + } else { + return Collections.emptyList(); + } + }); + try { + return list.get(); + } catch (InterruptedException | ExecutionException e) { + throw new RuntimeException(e); + } + } + + @Override + public ListenableFuture cleanUp(TenantId tenantId) { + return queueExecutor.submit(() -> { + data.remove(tenantId); + return null; + }); + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgQueue.java deleted file mode 100644 index 2825331623..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/InMemoryMsgQueue.java +++ /dev/null @@ -1,114 +0,0 @@ -/** - * Copyright © 2016-2018 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.dao.sql.queue; - -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningExecutorService; -import com.google.common.util.concurrent.MoreExecutors; -import lombok.Getter; -import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.queue.MsgQueue; -import org.thingsboard.server.dao.util.SqlDao; - -import javax.annotation.PostConstruct; -import javax.annotation.PreDestroy; -import java.util.*; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.Executors; -import java.util.concurrent.atomic.AtomicLong; - -/** - * Created by ashvayka on 27.04.18. - */ -@Component -//@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "memory", matchIfMissing = true) -@Slf4j -@SqlDao -public class InMemoryMsgQueue implements MsgQueue { - - @Value("${rule.queue.max_size}") - @Getter - private long maxSize; - - private ListeningExecutorService queueExecutor; - private AtomicLong pendingMsgCount = new AtomicLong(); - private Map> data = new HashMap<>(); - - @PostConstruct - public void init() { - // Should be always single threaded due to absence of locks. - queueExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); - } - - @PreDestroy - public void stop() { - if (queueExecutor == null) { - queueExecutor.shutdownNow(); - } - } - - @Override - public ListenableFuture put(TbMsg msg, UUID nodeId, long clusterPartition) { - if (pendingMsgCount.incrementAndGet() < maxSize) { - return queueExecutor.submit(() -> { - data.computeIfAbsent(new InMemoryMsgKey(nodeId, clusterPartition), key -> new HashMap<>()).put(msg.getId(), msg); - return null; - }); - } else { - pendingMsgCount.decrementAndGet(); - return Futures.immediateFailedFuture(new RuntimeException("Message queue is full!")); - } - } - - @Override - public ListenableFuture ack(TbMsg msg, UUID nodeId, long clusterPartition) { - return queueExecutor.submit(() -> { - InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); - Map map = data.get(key); - if (map != null) { - if (map.remove(msg.getId()) != null) { - pendingMsgCount.decrementAndGet(); - } - if (map.isEmpty()) { - data.remove(key); - } - } - return null; - }); - } - - @Override - public Iterable findUnprocessed(UUID nodeId, long clusterPartition) { - ListenableFuture> list = queueExecutor.submit(() -> { - InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); - Map map = data.get(key); - if (map != null) { - return new ArrayList<>(map.values()); - } else { - return Collections.emptyList(); - } - }); - try { - return list.get(); - } catch (InterruptedException | ExecutionException e) { - throw new RuntimeException(e); - } - } -} diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/QueuePartitionerTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java similarity index 92% rename from dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/QueuePartitionerTest.java rename to dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java index 2d3b61fec1..3de9542319 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/QueuePartitionerTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra; +package org.thingsboard.server.dao.queue.db.nosql; import org.junit.Before; @@ -21,7 +21,8 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; import org.mockito.runners.MockitoJUnitRunner; -import org.thingsboard.server.dao.service.queue.cassandra.repository.ProcessedPartitionRepository; +import org.thingsboard.server.dao.queue.db.nosql.QueuePartitioner; +import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; import java.time.Clock; import java.time.Instant; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/UnprocessedMsgFilterTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java similarity index 89% rename from dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/UnprocessedMsgFilterTest.java rename to dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java index caccc85932..39a432d244 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/UnprocessedMsgFilterTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java @@ -13,11 +13,13 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra; +package org.thingsboard.server.dao.queue.db.nosql; import com.google.common.collect.Lists; import org.junit.Test; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.dao.queue.db.nosql.MsgAck; +import org.thingsboard.server.dao.queue.db.nosql.UnprocessedMsgFilter; import java.util.Collection; import java.util.List; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraAckRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java similarity index 94% rename from dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraAckRepositoryTest.java rename to dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java index cff4dc9d0f..f2b8c88f01 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraAckRepositoryTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java @@ -13,18 +13,17 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; +package org.thingsboard.server.dao.queue.db.nosql.repository; import com.datastax.driver.core.utils.UUIDs; import com.google.common.collect.Lists; import com.google.common.util.concurrent.ListenableFuture; -import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.server.dao.service.AbstractServiceTest; import org.thingsboard.server.dao.service.DaoNoSqlTest; -import org.thingsboard.server.dao.service.queue.cassandra.MsgAck; +import org.thingsboard.server.dao.queue.db.nosql.MsgAck; import java.util.List; import java.util.UUID; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraMsgRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java similarity index 97% rename from dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraMsgRepositoryTest.java rename to dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java index fa286aaea2..f31db877ef 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraMsgRepositoryTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java @@ -13,13 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; +package org.thingsboard.server.dao.queue.db.nosql.repository; //import static org.junit.jupiter.api.Assertions.*; import com.datastax.driver.core.utils.UUIDs; import com.google.common.util.concurrent.ListenableFuture; -import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.util.ReflectionTestUtils; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraProcessedPartitionRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java similarity index 97% rename from dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraProcessedPartitionRepositoryTest.java rename to dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java index 2ae810a221..1ad053c49e 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/queue/cassandra/repository/impl/CassandraProcessedPartitionRepositoryTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java @@ -13,12 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; +package org.thingsboard.server.dao.queue.db.nosql.repository; import com.datastax.driver.core.utils.UUIDs; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.util.ReflectionTestUtils;