From 02e2c14fc959156cd2a8b406f33c6af3e875c7f1 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 10:18:06 +0300 Subject: [PATCH] Cleanup --- .../src/main/resources/thingsboard.yml | 2 +- .../server/dao/queue/db/MsgAck.java | 32 ------- .../dao/queue/db/UnprocessedMsgFilter.java | 35 -------- .../dao/queue/db/nosql/CassandraMsgQueue.java | 87 ------------------- .../dao/queue/db/nosql/QueuePartitioner.java | 86 ------------------ .../repository/CassandraAckRepository.java | 68 --------------- .../repository/CassandraMsgRepository.java | 67 -------------- ...CassandraProcessedPartitionRepository.java | 64 -------------- .../queue/db/repository/AckRepository.java | 29 ------- .../queue/db/repository/MsgRepository.java | 30 ------- .../ProcessedPartitionRepository.java | 29 ------- .../server/dao/queue/db/sql/SqlMsgQueue.java | 20 ----- .../queue/db/nosql/QueuePartitionerTest.java | 81 ----------------- .../db/nosql/UnprocessedMsgFilterTest.java | 47 ---------- .../CassandraAckRepositoryTest.java | 82 ----------------- .../CassandraMsgRepositoryTest.java | 87 ------------------- ...andraProcessedPartitionRepositoryTest.java | 83 ------------------ 17 files changed, 1 insertion(+), 928 deletions(-) delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index a10ef7d7c0..188291a4e2 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -257,7 +257,7 @@ actors: # Errors for particular actor are persisted once per specified amount of milliseconds error_persist_frequency: "${ACTORS_RULE_NODE_ERROR_FREQUENCY:3000}" queue: - # Message queue type (memory or db) + # Message queue type type: "${ACTORS_RULE_QUEUE_TYPE:memory}" # Message queue maximum size (per tenant) max_size: "${ACTORS_RULE_QUEUE_MAX_SIZE:100}" diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java deleted file mode 100644 index a1b039a1a7..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java +++ /dev/null @@ -1,32 +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.queue.db; - -import lombok.Data; -import lombok.EqualsAndHashCode; - -import java.util.UUID; - -@Data -@EqualsAndHashCode -public class MsgAck { - - private final UUID msgId; - private final UUID nodeId; - private final long clusteredPartition; - private final long tsPartition; - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java deleted file mode 100644 index 66eaa6d46b..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java +++ /dev/null @@ -1,35 +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.queue.db; - -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.queue.db.MsgAck; - -import java.util.Collection; -import java.util.List; -import java.util.Set; -import java.util.UUID; -import java.util.stream.Collectors; - -@Component -public class UnprocessedMsgFilter { - - public Collection filter(List msgs, List acks) { - Set processedIds = acks.stream().map(MsgAck::getMsgId).collect(Collectors.toSet()); - return msgs.stream().filter(i -> !processedIds.contains(i.getId())).collect(Collectors.toList()); - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java deleted file mode 100644 index 9cc87b47b6..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java +++ /dev/null @@ -1,87 +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.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.queue.db.MsgAck; -import org.thingsboard.server.dao.queue.db.UnprocessedMsgFilter; -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 = "actors.rule.queue", value = "type", havingValue = "db") -@Slf4j -@NoSqlDao -public class CassandraMsgQueue implements MsgQueue { - - @Autowired - private MsgRepository msgRepository; - @Autowired - private AckRepository ackRepository; - @Autowired - private UnprocessedMsgFilter unprocessedMsgFilter; - @Autowired - private QueuePartitioner queuePartitioner; - - @Override - 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(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(TenantId tenantId, UUID nodeId, long clusterPartition) { - List unprocessedMsgs = Lists.newArrayList(); - for (Long tsPartition : queuePartitioner.findUnprocessedPartitions(nodeId, clusterPartition)) { - List msgs = msgRepository.findMsgs(nodeId, clusterPartition, tsPartition); - List acks = ackRepository.findAcks(nodeId, clusterPartition, tsPartition); - unprocessedMsgs.addAll(unprocessedMsgFilter.filter(msgs, acks)); - } - 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/queue/db/nosql/QueuePartitioner.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java deleted file mode 100644 index 6076d93e9f..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java +++ /dev/null @@ -1,86 +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.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.queue.db.repository.ProcessedPartitionRepository; -import org.thingsboard.server.dao.timeseries.TsPartitionDate; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.time.Clock; -import java.time.Instant; -import java.time.LocalDateTime; -import java.time.ZoneOffset; -import java.util.List; -import java.util.Optional; -import java.util.UUID; -import java.util.concurrent.TimeUnit; - -@Component -@Slf4j -@NoSqlDao -public class QueuePartitioner { - - private final TsPartitionDate tsFormat; - private ProcessedPartitionRepository processedPartitionRepository; - private Clock clock = Clock.systemUTC(); - - public QueuePartitioner(@Value("${cassandra.queue.partitioning}") String partitioning, - ProcessedPartitionRepository processedPartitionRepository) { - this.processedPartitionRepository = processedPartitionRepository; - Optional partition = TsPartitionDate.parse(partitioning); - if (partition.isPresent()) { - tsFormat = partition.get(); - } else { - log.warn("Incorrect configuration of partitioning {}", partitioning); - throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); - } - } - - public long getPartition(long ts) { - //TODO: use TsPartitionDate.truncateTo? - LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); - return tsFormat.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli(); - } - - public List findUnprocessedPartitions(UUID nodeId, long clusteredHash) { - Optional lastPartitionOption = processedPartitionRepository.findLastProcessedPartition(nodeId, clusteredHash); - long lastPartition = lastPartitionOption.orElse(System.currentTimeMillis() - TimeUnit.DAYS.toMillis(7)); - List unprocessedPartitions = Lists.newArrayList(); - - LocalDateTime current = LocalDateTime.ofInstant(Instant.ofEpochMilli(lastPartition), ZoneOffset.UTC); - LocalDateTime end = LocalDateTime.ofInstant(Instant.now(clock), ZoneOffset.UTC) - .plus(1L, tsFormat.getTruncateUnit()); - - while (current.isBefore(end)) { - current = current.plus(1L, tsFormat.getTruncateUnit()); - unprocessedPartitions.add(tsFormat.truncatedTo(current).toInstant(ZoneOffset.UTC).toEpochMilli()); - } - - return unprocessedPartitions; - } - - public void setClock(Clock clock) { - this.clock = clock; - } - - public void checkProcessedPartitions() { - //todo-vp: we need to implement this - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java deleted file mode 100644 index 6c59c5997c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java +++ /dev/null @@ -1,68 +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.queue.db.nosql.repository; - -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.PreparedStatement; -import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Row; -import com.google.common.base.Function; -import com.google.common.util.concurrent.Futures; -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.queue.db.MsgAck; -import org.thingsboard.server.dao.queue.db.repository.AckRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.ArrayList; -import java.util.List; -import java.util.UUID; - -@Component -@NoSqlDao -public class CassandraAckRepository extends CassandraAbstractDao implements AckRepository { - - @Value("${cassandra.queue.ack.ttl}") - private int ackQueueTtl; - - @Override - public ListenableFuture ack(MsgAck msgAck) { - String insert = "INSERT INTO msg_ack_queue (node_id, cluster_partition, ts_partition, msg_id) VALUES (?, ?, ?, ?) USING TTL ?"; - PreparedStatement statement = prepare(insert); - BoundStatement boundStatement = statement.bind(msgAck.getNodeId(), msgAck.getClusteredPartition(), - msgAck.getTsPartition(), msgAck.getMsgId(), ackQueueTtl); - ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); - return Futures.transform(resultSetFuture, (Function) input -> null); - } - - @Override - public List findAcks(UUID nodeId, long clusterPartition, long tsPartition) { - String select = "SELECT msg_id FROM msg_ack_queue WHERE " + - "node_id = ? AND cluster_partition = ? AND ts_partition = ?"; - PreparedStatement statement = prepare(select); - BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition); - ResultSet rows = executeRead(boundStatement); - List msgs = new ArrayList<>(); - for (Row row : rows) { - msgs.add(new MsgAck(row.getUUID("msg_id"), nodeId, clusterPartition, tsPartition)); - } - return msgs; - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java deleted file mode 100644 index a90b71d99c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java +++ /dev/null @@ -1,67 +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.queue.db.nosql.repository; - -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.PreparedStatement; -import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Row; -import com.google.common.base.Function; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -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.queue.db.repository.MsgRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.ArrayList; -import java.util.List; -import java.util.UUID; - -@Component -@NoSqlDao -public class CassandraMsgRepository extends CassandraAbstractDao implements MsgRepository { - - @Value("${cassandra.queue.msg.ttl}") - private int msqQueueTtl; - - @Override - public ListenableFuture save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs) { - String insert = "INSERT INTO msg_queue (node_id, cluster_partition, ts_partition, ts, msg) VALUES (?, ?, ?, ?, ?) USING TTL ?"; - PreparedStatement statement = prepare(insert); - BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition, msgTs, TbMsg.toBytes(msg), msqQueueTtl); - ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); - return Futures.transform(resultSetFuture, (Function) input -> null); - } - - @Override - public List findMsgs(UUID nodeId, long clusterPartition, long tsPartition) { - String select = "SELECT node_id, cluster_partition, ts_partition, ts, msg FROM msg_queue WHERE " + - "node_id = ? AND cluster_partition = ? AND ts_partition = ?"; - PreparedStatement statement = prepare(select); - BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition); - ResultSet rows = executeRead(boundStatement); - List msgs = new ArrayList<>(); - for (Row row : rows) { - msgs.add(TbMsg.fromBytes(row.getBytes("msg"))); - } - return msgs; - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java deleted file mode 100644 index 831c6fdf4b..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java +++ /dev/null @@ -1,64 +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.queue.db.nosql.repository; - -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.PreparedStatement; -import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Row; -import com.google.common.base.Function; -import com.google.common.util.concurrent.Futures; -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.queue.db.repository.ProcessedPartitionRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.Optional; -import java.util.UUID; - -@Component -@NoSqlDao -public class CassandraProcessedPartitionRepository extends CassandraAbstractDao implements ProcessedPartitionRepository { - - @Value("${cassandra.queue.partitions.ttl}") - private int partitionsTtl; - - @Override - public ListenableFuture partitionProcessed(UUID nodeId, long clusterPartition, long tsPartition) { - String insert = "INSERT INTO processed_msg_partitions (node_id, cluster_partition, ts_partition) VALUES (?, ?, ?) USING TTL ?"; - PreparedStatement prepared = prepare(insert); - BoundStatement boundStatement = prepared.bind(nodeId, clusterPartition, tsPartition, partitionsTtl); - ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); - return Futures.transform(resultSetFuture, (Function) input -> null); - } - - @Override - public Optional findLastProcessedPartition(UUID nodeId, long clusteredHash) { - String select = "SELECT ts_partition FROM processed_msg_partitions WHERE " + - "node_id = ? AND cluster_partition = ?"; - PreparedStatement prepared = prepare(select); - BoundStatement boundStatement = prepared.bind(nodeId, clusteredHash); - Row row = executeRead(boundStatement).one(); - if (row == null) { - return Optional.empty(); - } - - return Optional.of(row.getLong("ts_partition")); - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java deleted file mode 100644 index 6fbd2da57e..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java +++ /dev/null @@ -1,29 +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.queue.db.repository; - -import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.dao.queue.db.MsgAck; - -import java.util.List; -import java.util.UUID; - -public interface AckRepository { - - ListenableFuture ack(MsgAck msgAck); - - List findAcks(UUID nodeId, long clusterPartition, long tsPartition); -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java deleted file mode 100644 index 0ca6900fc9..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java +++ /dev/null @@ -1,30 +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.queue.db.repository; - -import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.common.msg.TbMsg; - -import java.util.List; -import java.util.UUID; - -public interface MsgRepository { - - ListenableFuture save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs); - - List findMsgs(UUID nodeId, long clusterPartition, long tsPartition); - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java deleted file mode 100644 index b11fc6cc1c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java +++ /dev/null @@ -1,29 +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.queue.db.repository; - -import com.google.common.util.concurrent.ListenableFuture; - -import java.util.Optional; -import java.util.UUID; - -public interface ProcessedPartitionRepository { - - ListenableFuture partitionProcessed(UUID nodeId, long clusteredHash, long partition); - - Optional findLastProcessedPartition(UUID nodeId, long clusteredHash); - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java deleted file mode 100644 index f9dda43a36..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java +++ /dev/null @@ -1,20 +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.queue.db.sql; - -//@todo-vp: implement -public class SqlMsgQueue { -} diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java deleted file mode 100644 index e76ca2f946..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java +++ /dev/null @@ -1,81 +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.queue.db.nosql; - - -import org.junit.Before; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.mockito.Mock; -import org.mockito.runners.MockitoJUnitRunner; -import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; - -import java.time.Clock; -import java.time.Instant; -import java.time.ZoneOffset; -import java.time.temporal.ChronoUnit; -import java.util.List; -import java.util.Optional; -import java.util.UUID; - -import static org.junit.Assert.assertEquals; -import static org.mockito.Mockito.when; - -@RunWith(MockitoJUnitRunner.class) -public class QueuePartitionerTest { - - private QueuePartitioner queuePartitioner; - - @Mock - private ProcessedPartitionRepository partitionRepo; - - private Instant startInstant; - private Instant endInstant; - - @Before - public void init() { - queuePartitioner = new QueuePartitioner("MINUTES", partitionRepo); - startInstant = Instant.now(); - endInstant = startInstant.plus(2, ChronoUnit.MINUTES); - queuePartitioner.setClock(Clock.fixed(endInstant, ZoneOffset.UTC)); - } - - @Test - public void partitionCalculated() { - long time = 1519390191425L; - long partition = queuePartitioner.getPartition(time); - assertEquals(1519390140000L, partition); - } - - @Test - public void unprocessedPartitionsReturned() { - UUID nodeId = UUID.randomUUID(); - long clusteredHash = 101L; - when(partitionRepo.findLastProcessedPartition(nodeId, clusteredHash)).thenReturn(Optional.of(startInstant.toEpochMilli())); - List actual = queuePartitioner.findUnprocessedPartitions(nodeId, clusteredHash); - assertEquals(3, actual.size()); - } - - @Test - public void defaultShiftUsedIfNoPartitionWasProcessed() { - UUID nodeId = UUID.randomUUID(); - long clusteredHash = 101L; - when(partitionRepo.findLastProcessedPartition(nodeId, clusteredHash)).thenReturn(Optional.empty()); - List actual = queuePartitioner.findUnprocessedPartitions(nodeId, clusteredHash); - assertEquals(10083, actual.size()); - } - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java deleted file mode 100644 index fd9bf21164..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java +++ /dev/null @@ -1,47 +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.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.MsgAck; -import org.thingsboard.server.dao.queue.db.UnprocessedMsgFilter; - -import java.util.Collection; -import java.util.List; -import java.util.UUID; - -import static org.junit.Assert.assertEquals; - -public class UnprocessedMsgFilterTest { - - private UnprocessedMsgFilter msgFilter = new UnprocessedMsgFilter(); - - @Test - public void acknowledgedMsgsAreFilteredOut() { - UUID id1 = UUID.randomUUID(); - UUID id2 = UUID.randomUUID(); - TbMsg msg1 = new TbMsg(id1, "T", null, null, null, null, null, null, 0L); - TbMsg msg2 = new TbMsg(id2, "T", null, null, null, null, null, null, 0L); - List msgs = Lists.newArrayList(msg1, msg2); - List acks = Lists.newArrayList(new MsgAck(id2, UUID.randomUUID(), 1L, 1L)); - Collection actual = msgFilter.filter(msgs, acks); - assertEquals(1, actual.size()); - assertEquals(msg1, actual.iterator().next()); - } - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java deleted file mode 100644 index b2f38dc539..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java +++ /dev/null @@ -1,82 +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.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.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.queue.db.MsgAck; - -import java.util.List; -import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; - -@DaoNoSqlTest -public class CassandraAckRepositoryTest extends AbstractServiceTest { - - @Autowired - private CassandraAckRepository ackRepository; - - @Test - public void acksInPartitionCouldBeFound() { - UUID nodeId = UUID.fromString("055eee50-1883-11e8-b380-65b5d5335ba9"); - - List extectedAcks = Lists.newArrayList( - new MsgAck(UUID.fromString("bebaeb60-1888-11e8-bf21-65b5d5335ba9"), nodeId, 101L, 300L), - new MsgAck(UUID.fromString("12baeb60-1888-11e8-bf21-65b5d5335ba9"), nodeId, 101L, 300L) - ); - - List actualAcks = ackRepository.findAcks(nodeId, 101L, 300L); - assertEquals(extectedAcks, actualAcks); - } - - @Test - public void ackCanBeSavedAndRead() throws ExecutionException, InterruptedException { - UUID msgId = UUIDs.timeBased(); - UUID nodeId = UUIDs.timeBased(); - MsgAck ack = new MsgAck(msgId, nodeId, 10L, 20L); - ListenableFuture future = ackRepository.ack(ack); - future.get(); - List actualAcks = ackRepository.findAcks(nodeId, 10L, 20L); - assertEquals(1, actualAcks.size()); - assertEquals(ack, actualAcks.get(0)); - } - - @Test - public void expiredAcksAreNotReturned() throws ExecutionException, InterruptedException { - ReflectionTestUtils.setField(ackRepository, "ackQueueTtl", 1); - UUID msgId = UUIDs.timeBased(); - UUID nodeId = UUIDs.timeBased(); - MsgAck ack = new MsgAck(msgId, nodeId, 30L, 40L); - ListenableFuture future = ackRepository.ack(ack); - future.get(); - List actualAcks = ackRepository.findAcks(nodeId, 30L, 40L); - assertEquals(1, actualAcks.size()); - TimeUnit.SECONDS.sleep(2); - assertTrue(ackRepository.findAcks(nodeId, 30L, 40L).isEmpty()); - } - - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java deleted file mode 100644 index f31db877ef..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java +++ /dev/null @@ -1,87 +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.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.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.test.util.ReflectionTestUtils; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.RuleChainId; -import org.thingsboard.server.common.data.id.RuleNodeId; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.TbMsgDataType; -import org.thingsboard.server.common.msg.TbMsgMetaData; -import org.thingsboard.server.dao.service.AbstractServiceTest; -import org.thingsboard.server.dao.service.DaoNoSqlTest; - -import java.util.List; -import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; - -@DaoNoSqlTest -public class CassandraMsgRepositoryTest extends AbstractServiceTest { - - @Autowired - private CassandraMsgRepository msgRepository; - - @Test - public void msgCanBeSavedAndRead() throws ExecutionException, InterruptedException { - TbMsg msg = new TbMsg(UUIDs.timeBased(), "type", new DeviceId(UUIDs.timeBased()), null, TbMsgDataType.JSON, "0000", - new RuleChainId(UUIDs.timeBased()), new RuleNodeId(UUIDs.timeBased()), 0L); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = msgRepository.save(msg, nodeId, 1L, 1L, 1L); - future.get(); - List msgs = msgRepository.findMsgs(nodeId, 1L, 1L); - assertEquals(1, msgs.size()); - } - - @Test - public void expiredMsgsAreNotReturned() throws ExecutionException, InterruptedException { - ReflectionTestUtils.setField(msgRepository, "msqQueueTtl", 1); - TbMsg msg = new TbMsg(UUIDs.timeBased(), "type", new DeviceId(UUIDs.timeBased()), null, TbMsgDataType.JSON, "0000", - new RuleChainId(UUIDs.timeBased()), new RuleNodeId(UUIDs.timeBased()), 0L); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = msgRepository.save(msg, nodeId, 2L, 2L, 2L); - future.get(); - TimeUnit.SECONDS.sleep(2); - assertTrue(msgRepository.findMsgs(nodeId, 2L, 2L).isEmpty()); - } - - @Test - public void protoBufConverterWorkAsExpected() throws ExecutionException, InterruptedException { - TbMsgMetaData metaData = new TbMsgMetaData(); - metaData.putValue("key", "value"); - String dataStr = "someContent"; - TbMsg msg = new TbMsg(UUIDs.timeBased(), "type", new DeviceId(UUIDs.timeBased()), metaData, TbMsgDataType.JSON, dataStr, - new RuleChainId(UUIDs.timeBased()), new RuleNodeId(UUIDs.timeBased()), 0L); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = msgRepository.save(msg, nodeId, 1L, 1L, 1L); - future.get(); - List msgs = msgRepository.findMsgs(nodeId, 1L, 1L); - assertEquals(1, msgs.size()); - assertEquals(msg, msgs.get(0)); - } - - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java deleted file mode 100644 index a76fd965ca..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java +++ /dev/null @@ -1,83 +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.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.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 java.util.List; -import java.util.Optional; -import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertTrue; - -@DaoNoSqlTest -public class CassandraProcessedPartitionRepositoryTest extends AbstractServiceTest { - - @Autowired - private CassandraProcessedPartitionRepository partitionRepository; - - @Test - public void lastProcessedPartitionCouldBeFound() { - UUID nodeId = UUID.fromString("055eee50-1883-11e8-b380-65b5d5335ba9"); - Optional lastProcessedPartition = partitionRepository.findLastProcessedPartition(nodeId, 101L); - assertTrue(lastProcessedPartition.isPresent()); - assertEquals((Long) 777L, lastProcessedPartition.get()); - } - - @Test - public void highestProcessedPartitionReturned() throws ExecutionException, InterruptedException { - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future1 = partitionRepository.partitionProcessed(nodeId, 303L, 100L); - ListenableFuture future2 = partitionRepository.partitionProcessed(nodeId, 303L, 200L); - ListenableFuture future3 = partitionRepository.partitionProcessed(nodeId, 303L, 10L); - ListenableFuture> allFutures = Futures.allAsList(future1, future2, future3); - allFutures.get(); - Optional actual = partitionRepository.findLastProcessedPartition(nodeId, 303L); - assertTrue(actual.isPresent()); - assertEquals((Long) 200L, actual.get()); - } - - @Test - public void expiredPartitionsAreNotReturned() throws ExecutionException, InterruptedException { - ReflectionTestUtils.setField(partitionRepository, "partitionsTtl", 1); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = partitionRepository.partitionProcessed(nodeId, 404L, 10L); - future.get(); - Optional actual = partitionRepository.findLastProcessedPartition(nodeId, 404L); - assertEquals((Long) 10L, actual.get()); - TimeUnit.SECONDS.sleep(2); - assertFalse(partitionRepository.findLastProcessedPartition(nodeId, 404L).isPresent()); - } - - @Test - public void ifNoPartitionsWereProcessedEmptyResultReturned() { - UUID nodeId = UUIDs.timeBased(); - Optional actual = partitionRepository.findLastProcessedPartition(nodeId, 505L); - assertFalse(actual.isPresent()); - } - -} \ No newline at end of file