39 changed files with 165 additions and 995 deletions
@ -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; |
|
||||
|
|
||||
} |
|
||||
@ -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<TbMsg> filter(List<TbMsg> msgs, List<MsgAck> acks) { |
|
||||
Set<UUID> processedIds = acks.stream().map(MsgAck::getMsgId).collect(Collectors.toSet()); |
|
||||
return msgs.stream().filter(i -> !processedIds.contains(i.getId())).collect(Collectors.toList()); |
|
||||
} |
|
||||
} |
|
||||
@ -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 = "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<Void> 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<Void> 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<TbMsg> findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { |
|
||||
List<TbMsg> unprocessedMsgs = Lists.newArrayList(); |
|
||||
for (Long tsPartition : queuePartitioner.findUnprocessedPartitions(nodeId, clusterPartition)) { |
|
||||
List<TbMsg> msgs = msgRepository.findMsgs(nodeId, clusterPartition, tsPartition); |
|
||||
List<MsgAck> acks = ackRepository.findAcks(nodeId, clusterPartition, tsPartition); |
|
||||
unprocessedMsgs.addAll(unprocessedMsgFilter.filter(msgs, acks)); |
|
||||
} |
|
||||
return unprocessedMsgs; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> cleanUp(TenantId tenantId) { |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
|
|
||||
private long getMsgTime(TbMsg msg) { |
|
||||
return UUIDs.unixTimestamp(msg.getId()); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -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<TsPartitionDate> 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<Long> findUnprocessedPartitions(UUID nodeId, long clusteredHash) { |
|
||||
Optional<Long> lastPartitionOption = processedPartitionRepository.findLastProcessedPartition(nodeId, clusteredHash); |
|
||||
long lastPartition = lastPartitionOption.orElse(System.currentTimeMillis() - TimeUnit.DAYS.toMillis(7)); |
|
||||
List<Long> 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
|
|
||||
} |
|
||||
} |
|
||||
@ -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<Void> 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<ResultSet, Void>) input -> null); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public List<MsgAck> 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<MsgAck> msgs = new ArrayList<>(); |
|
||||
for (Row row : rows) { |
|
||||
msgs.add(new MsgAck(row.getUUID("msg_id"), nodeId, clusterPartition, tsPartition)); |
|
||||
} |
|
||||
return msgs; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -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<Void> 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<ResultSet, Void>) input -> null); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public List<TbMsg> 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<TbMsg> msgs = new ArrayList<>(); |
|
||||
for (Row row : rows) { |
|
||||
msgs.add(TbMsg.fromBytes(row.getBytes("msg"))); |
|
||||
} |
|
||||
return msgs; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -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<Void> 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<ResultSet, Void>) input -> null); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Optional<Long> 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")); |
|
||||
} |
|
||||
} |
|
||||
@ -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<Void> ack(MsgAck msgAck); |
|
||||
|
|
||||
List<MsgAck> findAcks(UUID nodeId, long clusterPartition, long tsPartition); |
|
||||
} |
|
||||
@ -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<Void> save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs); |
|
||||
|
|
||||
List<TbMsg> findMsgs(UUID nodeId, long clusterPartition, long tsPartition); |
|
||||
|
|
||||
} |
|
||||
@ -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<Void> partitionProcessed(UUID nodeId, long clusteredHash, long partition); |
|
||||
|
|
||||
Optional<Long> findLastProcessedPartition(UUID nodeId, long clusteredHash); |
|
||||
|
|
||||
} |
|
||||
@ -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 { |
|
||||
} |
|
||||
@ -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<Long> 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<Long> actual = queuePartitioner.findUnprocessedPartitions(nodeId, clusteredHash); |
|
||||
assertEquals(10083, actual.size()); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -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<TbMsg> msgs = Lists.newArrayList(msg1, msg2); |
|
||||
List<MsgAck> acks = Lists.newArrayList(new MsgAck(id2, UUID.randomUUID(), 1L, 1L)); |
|
||||
Collection<TbMsg> actual = msgFilter.filter(msgs, acks); |
|
||||
assertEquals(1, actual.size()); |
|
||||
assertEquals(msg1, actual.iterator().next()); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -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<MsgAck> 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<MsgAck> 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<Void> future = ackRepository.ack(ack); |
|
||||
future.get(); |
|
||||
List<MsgAck> 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<Void> future = ackRepository.ack(ack); |
|
||||
future.get(); |
|
||||
List<MsgAck> actualAcks = ackRepository.findAcks(nodeId, 30L, 40L); |
|
||||
assertEquals(1, actualAcks.size()); |
|
||||
TimeUnit.SECONDS.sleep(2); |
|
||||
assertTrue(ackRepository.findAcks(nodeId, 30L, 40L).isEmpty()); |
|
||||
} |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -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<Void> future = msgRepository.save(msg, nodeId, 1L, 1L, 1L); |
|
||||
future.get(); |
|
||||
List<TbMsg> 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<Void> 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<Void> future = msgRepository.save(msg, nodeId, 1L, 1L, 1L); |
|
||||
future.get(); |
|
||||
List<TbMsg> msgs = msgRepository.findMsgs(nodeId, 1L, 1L); |
|
||||
assertEquals(1, msgs.size()); |
|
||||
assertEquals(msg, msgs.get(0)); |
|
||||
} |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -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<Long> 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<Void> future1 = partitionRepository.partitionProcessed(nodeId, 303L, 100L); |
|
||||
ListenableFuture<Void> future2 = partitionRepository.partitionProcessed(nodeId, 303L, 200L); |
|
||||
ListenableFuture<Void> future3 = partitionRepository.partitionProcessed(nodeId, 303L, 10L); |
|
||||
ListenableFuture<List<Void>> allFutures = Futures.allAsList(future1, future2, future3); |
|
||||
allFutures.get(); |
|
||||
Optional<Long> 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<Void> future = partitionRepository.partitionProcessed(nodeId, 404L, 10L); |
|
||||
future.get(); |
|
||||
Optional<Long> 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<Long> actual = partitionRepository.findLastProcessedPartition(nodeId, 505L); |
|
||||
assertFalse(actual.isPresent()); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
Loading…
Reference in new issue