DATAREDIS-719 - Rework FutureResult and implementations for Jedis/Lettuce.
We now encapsulate deferred results for pipelining and transactions entirely within FutureResult and its subtypes. FutureResult accepts a Supplier<T> for default values, if operations return null and reports whether its result requires conversion. JedisResult and LettuceResult are now top-level classes and no longer inner classes of their connection factories. Original pull request: #289.
This commit is contained in:
committed by
Mark Paluch
parent
6ea78c77dd
commit
435ace8f55
@@ -15,8 +15,17 @@
|
||||
*/
|
||||
package org.springframework.data.redis.connection;
|
||||
|
||||
import java.util.*;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.Properties;
|
||||
import java.util.Queue;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
@@ -725,7 +734,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean mSet(Map<byte[], byte[]> tuple) {
|
||||
return delegate.mSet(tuple);
|
||||
return convertAndReturn(delegate.mSet(tuple), identityConverter);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -923,7 +932,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean set(byte[] key, byte[] value) {
|
||||
return delegate.set(key, value);
|
||||
return convertAndReturn(delegate.set(key, value), identityConverter);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -932,7 +941,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean set(byte[] key, byte[] value, Expiration expiration, SetOption option) {
|
||||
return delegate.set(key, value, expiration, option);
|
||||
return convertAndReturn(delegate.set(key, value, expiration, option), identityConverter);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -959,7 +968,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean setEx(byte[] key, long seconds, byte[] value) {
|
||||
return delegate.setEx(key, seconds, value);
|
||||
return convertAndReturn(delegate.setEx(key, seconds, value), identityConverter);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -968,7 +977,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean pSetEx(byte[] key, long milliseconds, byte[] value) {
|
||||
return delegate.pSetEx(key, milliseconds, value);
|
||||
return convertAndReturn(delegate.pSetEx(key, milliseconds, value), identityConverter);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -2106,7 +2115,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean mSetString(Map<String, String> tuple) {
|
||||
return delegate.mSet(serialize(tuple));
|
||||
return mSet(serialize(tuple));
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -2241,7 +2250,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean set(String key, String value) {
|
||||
return delegate.set(serialize(key), serialize(value));
|
||||
return set(serialize(key), serialize(value));
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -2268,7 +2277,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Boolean setEx(String key, long seconds, String value) {
|
||||
return delegate.setEx(serialize(key), seconds, serialize(value));
|
||||
return setEx(serialize(key), seconds, serialize(value));
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -3441,7 +3450,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Set<String> zRangeByLex(String key) {
|
||||
return this.zRangeByLex(key, Range.unbounded());
|
||||
return zRangeByLex(key, Range.unbounded());
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -3450,7 +3459,7 @@ public class DefaultStringRedisConnection implements StringRedisConnection, Deco
|
||||
*/
|
||||
@Override
|
||||
public Set<String> zRangeByLex(String key, Range range) {
|
||||
return this.zRangeByLex(key, range, null);
|
||||
return zRangeByLex(key, range, null);
|
||||
}
|
||||
|
||||
/*
|
||||
|
||||
@@ -15,6 +15,8 @@
|
||||
*/
|
||||
package org.springframework.data.redis.connection;
|
||||
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
@@ -28,23 +30,60 @@ import org.springframework.lang.Nullable;
|
||||
*/
|
||||
public abstract class FutureResult<T> {
|
||||
|
||||
protected T resultHolder;
|
||||
private T resultHolder;
|
||||
private final Supplier<?> defaultConversionResult;
|
||||
|
||||
protected boolean status = false;
|
||||
private boolean status = false;
|
||||
|
||||
@SuppressWarnings("rawtypes") //
|
||||
protected @Nullable Converter converter;
|
||||
|
||||
/**
|
||||
* Create new {@link FutureResult} for given object actually holding the result itself.
|
||||
*
|
||||
* @param resultHolder must not be {@literal null}.
|
||||
*/
|
||||
public FutureResult(T resultHolder) {
|
||||
this.resultHolder = resultHolder;
|
||||
this(resultHolder, val -> val);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create new {@link FutureResult} for given object actually holding the result itself and a converter capable of
|
||||
* transforming the result via {@link #convert(Object)}.
|
||||
*
|
||||
* @param resultHolder must not be {@literal null}.
|
||||
* @param converter can be {@literal null} and will be defaulted to an identity converter {@code value -> value} to
|
||||
* preserve the original value.
|
||||
*/
|
||||
@SuppressWarnings("rawtypes")
|
||||
public FutureResult(T resultHolder, Converter converter) {
|
||||
this.resultHolder = resultHolder;
|
||||
this.converter = converter;
|
||||
public FutureResult(T resultHolder, @Nullable Converter converter) {
|
||||
this(resultHolder, converter, () -> null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create new {@link FutureResult} for given object actually holding the result itself and a converter capable of
|
||||
* transforming the result via {@link #convert(Object)}.
|
||||
*
|
||||
* @param resultHolder must not be {@literal null}.
|
||||
* @param converter can be {@literal null} and will be defaulted to an identity converter {@code value -> value} to
|
||||
* preserve the original value.
|
||||
* @param defaultConversionResult must not be {@literal null}.
|
||||
* @since 2.1
|
||||
*/
|
||||
@SuppressWarnings("rawtypes")
|
||||
public FutureResult(T resultHolder, @Nullable Converter converter, Supplier<?> defaultConversionResult) {
|
||||
|
||||
this.resultHolder = resultHolder;
|
||||
this.converter = converter != null ? converter : val -> val;
|
||||
this.defaultConversionResult = defaultConversionResult;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get the object holding the actual result.
|
||||
*
|
||||
* @return never {@literal null}.
|
||||
* @since 1.1
|
||||
*/
|
||||
public T getResultHolder() {
|
||||
return resultHolder;
|
||||
}
|
||||
@@ -60,14 +99,18 @@ public abstract class FutureResult<T> {
|
||||
public Object convert(@Nullable Object result) {
|
||||
|
||||
if (result == null) {
|
||||
return null;
|
||||
return computeDefaultResult(result);
|
||||
}
|
||||
|
||||
return (converter != null) ? converter.convert(result) : result;
|
||||
return computeDefaultResult(converter.convert(result));
|
||||
}
|
||||
|
||||
@Nullable
|
||||
private Object computeDefaultResult(@Nullable Object source) {
|
||||
return source != null ? source : defaultConversionResult.get();
|
||||
}
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
@Nullable
|
||||
public Converter getConverter() {
|
||||
return converter;
|
||||
}
|
||||
@@ -93,4 +136,12 @@ public abstract class FutureResult<T> {
|
||||
*/
|
||||
@Nullable
|
||||
public abstract Object get();
|
||||
|
||||
/**
|
||||
* Indicate whether or not the actual result needs to be {@link #convert(Object) converted} before handing over.
|
||||
*
|
||||
* @return
|
||||
* @since 2.1
|
||||
*/
|
||||
public abstract boolean seeksConversion();
|
||||
}
|
||||
|
||||
@@ -71,7 +71,7 @@ public class TransactionResultConverter<T> implements Converter<List<Object>, Li
|
||||
: new RedisSystemException("Error reading future result.", source);
|
||||
}
|
||||
if (!(futureResult.isStatus())) {
|
||||
convertedResults.add(futureResult.convert(result));
|
||||
convertedResults.add(futureResult.seeksConversion() ? futureResult.convert(result) : result);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -29,6 +29,7 @@ import java.util.Collections;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Queue;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
@@ -39,6 +40,8 @@ import org.springframework.data.redis.RedisConnectionFailureException;
|
||||
import org.springframework.data.redis.RedisSystemException;
|
||||
import org.springframework.data.redis.connection.*;
|
||||
import org.springframework.data.redis.connection.convert.TransactionResultConverter;
|
||||
import org.springframework.data.redis.connection.jedis.JedisResult.JedisResultBuilder;
|
||||
import org.springframework.data.redis.connection.jedis.JedisResult.JedisStatusResult;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
@@ -75,43 +78,9 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
private final int dbIndex;
|
||||
private final String clientName;
|
||||
private boolean convertPipelineAndTxResults = true;
|
||||
private List<FutureResult<Response<?>>> pipelinedResults = new ArrayList<>();
|
||||
private List<JedisResult> pipelinedResults = new ArrayList<>();
|
||||
private Queue<FutureResult<Response<?>>> txResults = new LinkedList<>();
|
||||
|
||||
class JedisResult extends FutureResult<Response<?>> {
|
||||
|
||||
public <T> JedisResult(Response<T> resultHolder, Converter<T, ?> converter) {
|
||||
super(resultHolder, converter);
|
||||
}
|
||||
|
||||
public <T> JedisResult(Response<T> resultHolder) {
|
||||
super(resultHolder);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public Object get() {
|
||||
|
||||
Object raw = resultHolder.get();
|
||||
if (!convertPipelineAndTxResults || raw == null) {
|
||||
return raw;
|
||||
}
|
||||
|
||||
return converter != null ? converter.convert(raw) : raw;
|
||||
}
|
||||
}
|
||||
|
||||
private class JedisStatusResult extends JedisResult {
|
||||
public JedisStatusResult(Response<?> resultHolder) {
|
||||
super(resultHolder);
|
||||
setStatus(true);
|
||||
}
|
||||
|
||||
public <T> JedisStatusResult(Response<T> resultHolder, Converter<T, ?> converter) {
|
||||
super(resultHolder, converter);
|
||||
setStatus(true);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Constructs a new <code>JedisConnection</code> instance.
|
||||
*
|
||||
@@ -285,9 +254,9 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
Response<Object> result = JedisClientUtils
|
||||
.getResponse(isPipelined() ? getRequiredPipeline() : getRequiredTransaction());
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisResult(result));
|
||||
pipeline(newJedisResult(result));
|
||||
} else {
|
||||
transaction(new JedisResult(result));
|
||||
transaction(newJedisResult(result));
|
||||
}
|
||||
return null;
|
||||
}
|
||||
@@ -414,11 +383,13 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
List<Object> results = new ArrayList<>();
|
||||
getRequiredPipeline().sync();
|
||||
Exception cause = null;
|
||||
for (FutureResult<Response<?>> result : pipelinedResults) {
|
||||
for (JedisResult result : pipelinedResults) {
|
||||
try {
|
||||
|
||||
Object data = result.get();
|
||||
if (!convertPipelineAndTxResults || !(result.isStatus())) {
|
||||
results.add(data);
|
||||
|
||||
if (!result.isStatus()) {
|
||||
results.add(result.seeksConversion() ? result.convert(data) : data);
|
||||
}
|
||||
} catch (JedisDataException e) {
|
||||
DataAccessException dataAccessException = convertJedisAccessException(e);
|
||||
@@ -439,7 +410,7 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
return results;
|
||||
}
|
||||
|
||||
void pipeline(FutureResult<Response<?>> result) {
|
||||
void pipeline(JedisResult result) {
|
||||
if (isQueueing()) {
|
||||
transaction(result);
|
||||
} else {
|
||||
@@ -459,11 +430,11 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
public byte[] echo(byte[] message) {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisResult(getRequiredPipeline().echo(message)));
|
||||
pipeline(newJedisResult(getRequiredPipeline().echo(message)));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(new JedisResult(getRequiredTransaction().echo(message)));
|
||||
transaction(newJedisResult(getRequiredTransaction().echo(message)));
|
||||
return null;
|
||||
}
|
||||
return jedis.echo(message);
|
||||
@@ -480,11 +451,11 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
public String ping() {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisResult(getRequiredPipeline().ping()));
|
||||
pipeline(newJedisResult(getRequiredPipeline().ping()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(new JedisResult(getRequiredTransaction().ping()));
|
||||
transaction(newJedisResult(getRequiredTransaction().ping()));
|
||||
return null;
|
||||
}
|
||||
return jedis.ping();
|
||||
@@ -501,7 +472,7 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
public void discard() {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisStatusResult(getRequiredPipeline().discard()));
|
||||
pipeline(newStatusResult(getRequiredPipeline().discard()));
|
||||
return;
|
||||
}
|
||||
getRequiredTransaction().discard();
|
||||
@@ -521,7 +492,7 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
public List<Object> exec() {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisResult(getRequiredPipeline().exec(),
|
||||
pipeline(newJedisResult(getRequiredPipeline().exec(),
|
||||
new TransactionResultConverter<>(new LinkedList<>(txResults), JedisConverters.exceptionConverter())));
|
||||
return null;
|
||||
}
|
||||
@@ -529,8 +500,10 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
if (transaction == null) {
|
||||
throw new InvalidDataAccessApiUsageException("No ongoing transaction. Did you forget to call multi?");
|
||||
}
|
||||
|
||||
List<Object> results = transaction.exec();
|
||||
return convertPipelineAndTxResults && !CollectionUtils.isEmpty(results)
|
||||
|
||||
return !CollectionUtils.isEmpty(results)
|
||||
? new TransactionResultConverter<>(txResults, JedisConverters.exceptionConverter()).convert(results)
|
||||
: results;
|
||||
} catch (Exception ex) {
|
||||
@@ -578,19 +551,23 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
}
|
||||
|
||||
JedisResult newJedisResult(Response<?> response) {
|
||||
return new JedisResult(response);
|
||||
return JedisResultBuilder.forResponse(response).build();
|
||||
}
|
||||
|
||||
<T> JedisResult newJedisResult(Response<T> response, Converter<T, ?> converter) {
|
||||
return new JedisResult(response, converter);
|
||||
|
||||
return JedisResultBuilder.forResponse(response).mappedWith(converter)
|
||||
.convertPipelineAndTxResults(convertPipelineAndTxResults).build();
|
||||
}
|
||||
|
||||
<T> JedisResult newJedisResult(Response<T> response, Converter<T, ?> converter, Supplier<?> defaultValue) {
|
||||
|
||||
return JedisResultBuilder.forResponse(response).mappedWith(converter)
|
||||
.convertPipelineAndTxResults(convertPipelineAndTxResults).defaultNullTo(defaultValue).build();
|
||||
}
|
||||
|
||||
JedisStatusResult newStatusResult(Response<?> response) {
|
||||
return new JedisStatusResult(response);
|
||||
}
|
||||
|
||||
<T> JedisStatusResult newStatusResult(Response<T> response, Converter<T, ?> converter) {
|
||||
return new JedisStatusResult(response, converter);
|
||||
return JedisResultBuilder.forResponse(response).buildStatusResult();
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -621,11 +598,11 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
public void select(int dbIndex) {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisStatusResult(getRequiredPipeline().select(dbIndex)));
|
||||
pipeline(newStatusResult(getRequiredPipeline().select(dbIndex)));
|
||||
return;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(new JedisStatusResult(getRequiredTransaction().select(dbIndex)));
|
||||
transaction(newStatusResult(getRequiredTransaction().select(dbIndex)));
|
||||
return;
|
||||
}
|
||||
jedis.select(dbIndex);
|
||||
@@ -659,7 +636,7 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
try {
|
||||
for (byte[] key : keys) {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisStatusResult(getRequiredPipeline().watch(key)));
|
||||
pipeline(newStatusResult(getRequiredPipeline().watch(key)));
|
||||
} else {
|
||||
jedis.watch(key);
|
||||
}
|
||||
@@ -681,11 +658,11 @@ public class JedisConnection extends AbstractRedisConnection {
|
||||
public Long publish(byte[] channel, byte[] message) {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new JedisResult(getRequiredPipeline().publish(channel, message)));
|
||||
pipeline(newJedisResult(getRequiredPipeline().publish(channel, message)));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(new JedisResult(getRequiredTransaction().publish(channel, message)));
|
||||
transaction(newJedisResult(getRequiredTransaction().publish(channel, message)));
|
||||
return null;
|
||||
}
|
||||
return jedis.publish(channel, message);
|
||||
|
||||
@@ -32,7 +32,6 @@ import org.springframework.data.geo.Metric;
|
||||
import org.springframework.data.geo.Point;
|
||||
import org.springframework.data.redis.connection.RedisGeoCommands;
|
||||
import org.springframework.data.redis.connection.convert.ListConverter;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
@@ -63,9 +62,8 @@ class JedisGeoCommands implements RedisGeoCommands {
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(
|
||||
connection
|
||||
.newJedisResult(connection.getRequiredTransaction().geoadd(key, point.getX(), point.getY(), member)));
|
||||
transaction(connection
|
||||
.newJedisResult(connection.getRequiredTransaction().geoadd(key, point.getX(), point.getY(), member)));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -159,9 +157,8 @@ class JedisGeoCommands implements RedisGeoCommands {
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(
|
||||
connection.newJedisResult(connection.getRequiredTransaction().geodist(key, member1, member2),
|
||||
distanceConverter));
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().geodist(key, member1, member2),
|
||||
distanceConverter));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -194,8 +191,7 @@ class JedisGeoCommands implements RedisGeoCommands {
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newJedisResult(
|
||||
connection.getRequiredTransaction().geodist(key, member1, member2, geoUnit),
|
||||
distanceConverter));
|
||||
connection.getRequiredTransaction().geodist(key, member1, member2, geoUnit), distanceConverter));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -323,8 +319,7 @@ class JedisGeoCommands implements RedisGeoCommands {
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().georadius(key,
|
||||
within.getCenter().getX(),
|
||||
within.getCenter().getY(), within.getRadius().getValue(),
|
||||
within.getCenter().getX(), within.getCenter().getY(), within.getRadius().getValue(),
|
||||
JedisConverters.toGeoUnit(within.getRadius().getMetric()), geoRadiusParam), converter));
|
||||
return null;
|
||||
}
|
||||
@@ -396,10 +391,8 @@ class JedisGeoCommands implements RedisGeoCommands {
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newJedisResult(
|
||||
connection.getRequiredTransaction().georadiusByMember(key, member, radius.getValue(), geoUnit,
|
||||
geoRadiusParam),
|
||||
converter));
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().georadiusByMember(key, member,
|
||||
radius.getValue(), geoUnit, geoRadiusParam), converter));
|
||||
return null;
|
||||
}
|
||||
return converter
|
||||
|
||||
@@ -26,7 +26,6 @@ import java.util.Map.Entry;
|
||||
import java.util.Set;
|
||||
|
||||
import org.springframework.data.redis.connection.RedisHashCommands;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.KeyBoundCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -19,7 +19,6 @@ import lombok.NonNull;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
|
||||
import org.springframework.data.redis.connection.RedisHyperLogLogCommands;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
|
||||
@@ -28,7 +28,6 @@ import org.springframework.data.redis.connection.DataType;
|
||||
import org.springframework.data.redis.connection.RedisKeyCommands;
|
||||
import org.springframework.data.redis.connection.SortParameters;
|
||||
import org.springframework.data.redis.connection.convert.Converters;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.ScanCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
@@ -83,13 +82,11 @@ class JedisKeyCommands implements RedisKeyCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(
|
||||
connection.newJedisResult(connection.getRequiredPipeline().exists(keys)));
|
||||
pipeline(connection.newJedisResult(connection.getRequiredPipeline().exists(keys)));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection
|
||||
.newJedisResult(connection.getRequiredTransaction().exists(keys)));
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().exists(keys)));
|
||||
return null;
|
||||
}
|
||||
return connection.getJedis().exists(keys);
|
||||
|
||||
@@ -23,7 +23,6 @@ import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.data.redis.connection.RedisListCommands;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
/*
|
||||
* Copyright 2017 the original author or 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.springframework.data.redis.connection.jedis;
|
||||
|
||||
import redis.clients.jedis.Response;
|
||||
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.data.redis.connection.FutureResult;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
/**
|
||||
* Jedis specific {@link FutureResult} implementation. <br />
|
||||
*
|
||||
* @author Costin Leau
|
||||
* @author Jennifer Hickey
|
||||
* @author Christoph Strobl
|
||||
* @author Mark Paluch
|
||||
* @since 2.1
|
||||
*/
|
||||
class JedisResult<T, S> extends FutureResult<Response<?>> {
|
||||
|
||||
private final boolean convertPipelineAndTxResults;
|
||||
|
||||
<T> JedisResult(Response<T> resultHolder) {
|
||||
this(resultHolder, false, null);
|
||||
}
|
||||
|
||||
<T> JedisResult(Response<T> resultHolder, boolean convertPipelineAndTxResults, @Nullable Converter<T, ?> converter) {
|
||||
this(resultHolder, null, convertPipelineAndTxResults, converter);
|
||||
}
|
||||
|
||||
<T> JedisResult(Response<T> resultHolder, Supplier<S> defaultReturnValue, boolean convertPipelineAndTxResults,
|
||||
@Nullable Converter<T, ?> converter) {
|
||||
|
||||
super(resultHolder, converter, defaultReturnValue);
|
||||
this.convertPipelineAndTxResults = convertPipelineAndTxResults;
|
||||
}
|
||||
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.redis.connection.FutureResult#get()
|
||||
* @return
|
||||
*/
|
||||
@Nullable
|
||||
@Override
|
||||
public T get() {
|
||||
return (T) getResultHolder().get();
|
||||
}
|
||||
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.redis.connection.FutureResult#seeksConversion()
|
||||
* @return
|
||||
*/
|
||||
public boolean seeksConversion() {
|
||||
return convertPipelineAndTxResults && converter != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Jedis specific {@link FutureResult} implementation of a throw away status result.
|
||||
*/
|
||||
static class JedisStatusResult extends JedisResult {
|
||||
|
||||
<T> JedisStatusResult(Response<T> resultHolder, Converter<T, ?> converter) {
|
||||
|
||||
super(resultHolder, false, converter);
|
||||
setStatus(true);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Builder for constructing {@link JedisResult}.
|
||||
*
|
||||
* @param <T>
|
||||
* @param <S>
|
||||
* @since 2.1
|
||||
*/
|
||||
static class JedisResultBuilder<T, S> {
|
||||
|
||||
private final Response<T> response;
|
||||
private Converter<T, ?> converter;
|
||||
private boolean convertPipelineAndTxResults = false;
|
||||
private Supplier<?> nullValueDefault = () -> null;
|
||||
|
||||
JedisResultBuilder(Response<T> response) {
|
||||
|
||||
this.response = response;
|
||||
this.converter = (source) -> source;
|
||||
}
|
||||
|
||||
static <T> JedisResultBuilder<T, ?> forResponse(Response<T> response) {
|
||||
return new JedisResultBuilder<>(response);
|
||||
}
|
||||
|
||||
<S> JedisResultBuilder<T, S> mappedWith(Converter<T, S> converter) {
|
||||
|
||||
this.converter = converter;
|
||||
return (JedisResultBuilder<T, S>) this;
|
||||
}
|
||||
|
||||
<S> JedisResultBuilder<T, S> defaultNullTo(S value) {
|
||||
return (defaultNullTo(() -> value));
|
||||
}
|
||||
|
||||
<S> JedisResultBuilder<T, S> defaultNullTo(Supplier<S> value) {
|
||||
|
||||
this.nullValueDefault = value;
|
||||
return (JedisResultBuilder<T, S>) this;
|
||||
}
|
||||
|
||||
JedisResultBuilder<T, S> convertPipelineAndTxResults(boolean flag) {
|
||||
|
||||
convertPipelineAndTxResults = flag;
|
||||
return this;
|
||||
}
|
||||
|
||||
JedisResult<T, S> build() {
|
||||
return new JedisResult(response, nullValueDefault, convertPipelineAndTxResults, converter);
|
||||
}
|
||||
|
||||
JedisStatusResult buildStatusResult() {
|
||||
return new JedisStatusResult(response, converter);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -25,7 +25,6 @@ import org.springframework.data.redis.connection.RedisNode;
|
||||
import org.springframework.data.redis.connection.RedisServerCommands;
|
||||
import org.springframework.data.redis.connection.ReturnType;
|
||||
import org.springframework.data.redis.connection.convert.Converters;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.data.redis.core.types.RedisClientInfo;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -24,7 +24,6 @@ import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
import org.springframework.data.redis.connection.RedisSetCommands;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.KeyBoundCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -24,7 +24,6 @@ import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.springframework.data.redis.connection.RedisStringCommands;
|
||||
import org.springframework.data.redis.connection.convert.Converters;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.data.redis.core.types.Expiration;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
@@ -126,12 +125,12 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newStatusResult(connection.getRequiredPipeline().set(key, value),
|
||||
pipeline(connection.newJedisResult(connection.getRequiredPipeline().set(key, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newStatusResult(connection.getRequiredTransaction().set(key, value),
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().set(key, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
@@ -165,14 +164,14 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
|
||||
if (isPipelined()) {
|
||||
|
||||
pipeline(connection.newStatusResult(connection.getRequiredPipeline().set(key, value, nxxx),
|
||||
Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newJedisResult(connection.getRequiredPipeline().set(key, value, nxxx),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
|
||||
transaction(connection.newStatusResult(connection.getRequiredTransaction().set(key, value, nxxx),
|
||||
Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().set(key, value, nxxx),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -205,9 +204,9 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
"Expiration.expirationTime must be less than Integer.MAX_VALUE for pipeline in Jedis.");
|
||||
}
|
||||
|
||||
pipeline(connection.newStatusResult(
|
||||
pipeline(connection.newJedisResult(
|
||||
connection.getRequiredPipeline().set(key, value, nxxx, expx, (int) expiration.getExpirationTime()),
|
||||
Converters.stringToBooleanConverter()));
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
@@ -217,9 +216,9 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
"Expiration.expirationTime must be less than Integer.MAX_VALUE for transactions in Jedis.");
|
||||
}
|
||||
|
||||
transaction(connection.newStatusResult(
|
||||
transaction(connection.newJedisResult(
|
||||
connection.getRequiredTransaction().set(key, value, nxxx, expx, (int) expiration.getExpirationTime()),
|
||||
Converters.stringToBooleanConverter()));
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -277,13 +276,13 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newStatusResult(connection.getRequiredPipeline().setex(key, (int) seconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newJedisResult(connection.getRequiredPipeline().setex(key, (int) seconds, value),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newStatusResult(connection.getRequiredTransaction().setex(key, (int) seconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().setex(key, (int) seconds, value),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(connection.getJedis().setex(key, (int) seconds, value));
|
||||
@@ -304,13 +303,13 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newStatusResult(connection.getRequiredPipeline().psetex(key, milliseconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newJedisResult(connection.getRequiredPipeline().psetex(key, milliseconds, value),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newStatusResult(connection.getRequiredTransaction().psetex(key, milliseconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newJedisResult(connection.getRequiredTransaction().psetex(key, milliseconds, value),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(connection.getJedis().psetex(key, milliseconds, value));
|
||||
@@ -330,13 +329,13 @@ class JedisStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newStatusResult(connection.getRequiredPipeline().mset(JedisConverters.toByteArrays(tuples)),
|
||||
pipeline(connection.newJedisResult(connection.getRequiredPipeline().mset(JedisConverters.toByteArrays(tuples)),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(
|
||||
connection.newStatusResult(connection.getRequiredTransaction().mset(JedisConverters.toByteArrays(tuples)),
|
||||
connection.newJedisResult(connection.getRequiredTransaction().mset(JedisConverters.toByteArrays(tuples)),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -25,7 +25,6 @@ import java.nio.charset.StandardCharsets;
|
||||
import java.util.Set;
|
||||
|
||||
import org.springframework.data.redis.connection.RedisZSetCommands;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnection.JedisResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.KeyBoundCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -53,6 +53,7 @@ import java.util.Queue;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.springframework.beans.BeanUtils;
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
@@ -64,6 +65,10 @@ import org.springframework.data.redis.FallbackExceptionTranslationStrategy;
|
||||
import org.springframework.data.redis.connection.*;
|
||||
import org.springframework.data.redis.connection.convert.TransactionResultConverter;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnectionProvider.TargetAware;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceResultBuilder;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceStatusResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxStatusResult;
|
||||
import org.springframework.data.redis.core.RedisCommand;
|
||||
import org.springframework.data.redis.core.ScanOptions;
|
||||
import org.springframework.lang.Nullable;
|
||||
@@ -110,106 +115,46 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
/** flag indicating whether the connection needs to be dropped or not */
|
||||
private boolean convertPipelineAndTxResults = true;
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
class LettuceResult extends FutureResult<io.lettuce.core.protocol.RedisCommand<?, ?, ?>> {
|
||||
public <T> LettuceResult(Future<T> resultHolder, Converter<T, ?> converter) {
|
||||
super((io.lettuce.core.protocol.RedisCommand) resultHolder, converter);
|
||||
}
|
||||
|
||||
public LettuceResult(Future resultHolder) {
|
||||
super((io.lettuce.core.protocol.RedisCommand) resultHolder);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
public Object get() {
|
||||
try {
|
||||
if (convertPipelineAndTxResults && converter != null) {
|
||||
return converter.convert(resultHolder.getOutput().get());
|
||||
}
|
||||
return resultHolder.getOutput().get();
|
||||
} catch (Exception e) {
|
||||
throw LettuceConnection.this.convertLettuceAccessException(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
LettuceResult newLettuceResult(Future<?> resultHolder) {
|
||||
return new LettuceResult(resultHolder);
|
||||
return newLettuceResult(resultHolder, (val) -> val);
|
||||
}
|
||||
|
||||
<T> LettuceResult newLettuceResult(Future<T> resultHolder, Converter<T, ?> converter) {
|
||||
return new LettuceResult(resultHolder, converter);
|
||||
|
||||
return LettuceResultBuilder.forResponse(resultHolder).mappedWith(converter)
|
||||
.convertPipelineAndTxResults(convertPipelineAndTxResults).build();
|
||||
}
|
||||
|
||||
private class LettuceStatusResult extends LettuceResult {
|
||||
@SuppressWarnings("rawtypes")
|
||||
LettuceStatusResult(Future resultHolder) {
|
||||
super(resultHolder);
|
||||
setStatus(true);
|
||||
}
|
||||
<T> LettuceResult newLettuceResult(Future<T> resultHolder, Converter<T, ?> converter, Supplier<?> defaultValue) {
|
||||
|
||||
<T> LettuceStatusResult(Future<T> resultHolder, Converter<T, ?> converter) {
|
||||
super(resultHolder, converter);
|
||||
setStatus(true);
|
||||
}
|
||||
return LettuceResultBuilder.forResponse(resultHolder).mappedWith(converter)
|
||||
.convertPipelineAndTxResults(convertPipelineAndTxResults).defaultNullTo(defaultValue).build();
|
||||
}
|
||||
|
||||
LettuceStatusResult newLettuceStatusResult(Future<?> resultHolder) {
|
||||
return new LettuceStatusResult(resultHolder);
|
||||
}
|
||||
|
||||
<T> LettuceStatusResult newLettuceStatusResult(Future<T> resultHolder, Converter<T, ?> converter) {
|
||||
return new LettuceStatusResult(resultHolder, converter);
|
||||
}
|
||||
|
||||
class LettuceTxResult extends FutureResult<Object> {
|
||||
public LettuceTxResult(Object resultHolder, Converter<?, ?> converter) {
|
||||
super(resultHolder, converter);
|
||||
}
|
||||
|
||||
public LettuceTxResult(Object resultHolder) {
|
||||
super(resultHolder);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
public Object get() {
|
||||
if (convertPipelineAndTxResults && converter != null) {
|
||||
return converter.convert(resultHolder);
|
||||
}
|
||||
return resultHolder;
|
||||
}
|
||||
}
|
||||
|
||||
LettuceTxResult newLettuceTxResult(Object resultHolder) {
|
||||
return new LettuceTxResult(resultHolder);
|
||||
return newLettuceTxResult(resultHolder, (val) -> val);
|
||||
}
|
||||
|
||||
LettuceTxResult newLettuceTxResult(Object resultHolder, Converter<?, ?> converter) {
|
||||
return new LettuceTxResult(resultHolder, converter);
|
||||
<T> LettuceTxResult<T> newLettuceTxResult(Object resultHolder, Converter converter) {
|
||||
|
||||
return LettuceResultBuilder.forResponse(resultHolder).mappedWith(converter)
|
||||
.convertPipelineAndTxResults(convertPipelineAndTxResults).buildTxResult();
|
||||
}
|
||||
|
||||
private class LettuceTxStatusResult extends LettuceTxResult {
|
||||
LettuceTxStatusResult(Object resultHolder) {
|
||||
super(resultHolder);
|
||||
setStatus(true);
|
||||
}
|
||||
<T> LettuceTxResult<T> newLettuceTxResult(Object resultHolder, Converter converter, Supplier<?> defaultValue) {
|
||||
|
||||
LettuceTxStatusResult(Object resultHolder, Converter converter) {
|
||||
super(resultHolder, converter);
|
||||
setStatus(true);
|
||||
}
|
||||
return LettuceResultBuilder.forResponse(resultHolder).mappedWith(converter)
|
||||
.convertPipelineAndTxResults(convertPipelineAndTxResults).defaultNullTo(defaultValue).buildTxResult();
|
||||
}
|
||||
|
||||
LettuceTxStatusResult newLettuceTxStatusResult(Object resultHolder) {
|
||||
return new LettuceTxStatusResult(resultHolder);
|
||||
}
|
||||
|
||||
LettuceTxStatusResult newLettuceTxStatusResult(Object resultHolder, Converter<?, ?> converter) {
|
||||
return new LettuceTxStatusResult(resultHolder, converter);
|
||||
}
|
||||
|
||||
private class LettuceTransactionResultConverter<T> extends TransactionResultConverter<T> {
|
||||
public LettuceTransactionResultConverter(Queue<FutureResult<T>> txResults,
|
||||
Converter<Exception, DataAccessException> exceptionConverter) {
|
||||
@@ -440,7 +385,7 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
/**
|
||||
* 'Native' or 'raw' execution of the given command along-side the given arguments.
|
||||
*
|
||||
* @see RedisCommands#execute(String, byte[]...)
|
||||
* @see RedisConnection#execute(String, byte[]...)
|
||||
* @param command Command to execute
|
||||
* @param commandOutputTypeHint Type of Output to use, may be (may be {@literal null}).
|
||||
* @param args Possible command arguments (may be {@literal null})
|
||||
@@ -469,11 +414,11 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
Command cmd = new Command(commandType, expectedOutput, cmdArg);
|
||||
if (isPipelined()) {
|
||||
|
||||
pipeline(new LettuceResult(connectionImpl.dispatch(cmd.getType(), cmd.getOutput(), cmd.getArgs())));
|
||||
pipeline(newLettuceResult(connectionImpl.dispatch(cmd.getType(), cmd.getOutput(), cmd.getArgs())));
|
||||
return null;
|
||||
} else if (isQueueing()) {
|
||||
|
||||
transaction(new LettuceTxResult(connectionImpl.dispatch(cmd.getType(), cmd.getOutput(), cmd.getArgs())));
|
||||
transaction(newLettuceTxResult(connectionImpl.dispatch(cmd.getType(), cmd.getOutput(), cmd.getArgs())));
|
||||
return null;
|
||||
} else {
|
||||
return await(connectionImpl.dispatch(cmd.getType(), cmd.getOutput(), cmd.getArgs()));
|
||||
@@ -546,7 +491,7 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
if (isPipelined) {
|
||||
isPipelined = false;
|
||||
List<io.lettuce.core.protocol.RedisCommand<?, ?, ?>> futures = new ArrayList<>();
|
||||
for (LettuceResult result : ppline) {
|
||||
for (LettuceResult<?, ?> result : ppline) {
|
||||
futures.add(result.getResultHolder());
|
||||
}
|
||||
|
||||
@@ -559,7 +504,7 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
Exception problem = null;
|
||||
|
||||
if (done) {
|
||||
for (LettuceResult result : ppline) {
|
||||
for (LettuceResult<?, ?> result : ppline) {
|
||||
|
||||
if (result.getResultHolder().getOutput().hasError()) {
|
||||
|
||||
@@ -569,10 +514,10 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
problem = err;
|
||||
}
|
||||
results.add(err);
|
||||
} else if (!convertPipelineAndTxResults || !(result.isStatus())) {
|
||||
} else if (!result.isStatus()) {
|
||||
|
||||
try {
|
||||
results.add(result.get());
|
||||
results.add(result.seeksConversion() ? result.convert(result.get()) : result.get());
|
||||
} catch (DataAccessException e) {
|
||||
if (problem == null) {
|
||||
problem = e;
|
||||
@@ -609,7 +554,7 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
public byte[] echo(byte[] message) {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new LettuceResult(getAsyncConnection().echo(message)));
|
||||
pipeline(newLettuceResult(getAsyncConnection().echo(message)));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
@@ -626,7 +571,7 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
public String ping() {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new LettuceResult(getAsyncConnection().ping()));
|
||||
pipeline(newLettuceResult(getAsyncConnection().ping()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
@@ -658,7 +603,9 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
@Override
|
||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||
public List<Object> exec() {
|
||||
|
||||
isMulti = false;
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
RedisFuture<TransactionResult> exec = ((RedisAsyncCommands) getAsyncDedicatedConnection()).exec();
|
||||
@@ -666,8 +613,8 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
LettuceTransactionResultConverter resultConverter = new LettuceTransactionResultConverter(
|
||||
new LinkedList<>(txResults), LettuceConverters.exceptionConverter());
|
||||
|
||||
pipeline(new LettuceResult(exec,
|
||||
source -> resultConverter.convert(LettuceConverters.transactionResultUnwrapper().convert(source))));
|
||||
pipeline(newLettuceResult(exec, source -> resultConverter
|
||||
.convert(LettuceConverters.transactionResultUnwrapper().convert((TransactionResult) source))));
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -768,11 +715,11 @@ public class LettuceConnection extends AbstractRedisConnection {
|
||||
public Long publish(byte[] channel, byte[] message) {
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(new LettuceResult(getAsyncConnection().publish(channel, message)));
|
||||
pipeline(newLettuceResult(getAsyncConnection().publish(channel, message)));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(new LettuceTxResult(getConnection().publish(channel, message)));
|
||||
transaction(newLettuceTxResult(getConnection().publish(channel, message)));
|
||||
return null;
|
||||
}
|
||||
return getConnection().publish(channel, message);
|
||||
|
||||
@@ -40,8 +40,7 @@ import org.springframework.data.geo.Metric;
|
||||
import org.springframework.data.geo.Point;
|
||||
import org.springframework.data.redis.connection.RedisGeoCommands;
|
||||
import org.springframework.data.redis.connection.convert.ListConverter;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
|
||||
@@ -29,8 +29,7 @@ import java.util.Set;
|
||||
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisHashCommands;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.KeyBoundCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -24,8 +24,7 @@ import lombok.RequiredArgsConstructor;
|
||||
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisHyperLogLogCommands;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
|
||||
@@ -32,8 +32,7 @@ import org.springframework.data.redis.connection.DataType;
|
||||
import org.springframework.data.redis.connection.RedisKeyCommands;
|
||||
import org.springframework.data.redis.connection.SortParameters;
|
||||
import org.springframework.data.redis.connection.convert.Converters;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.ScanCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -24,8 +24,7 @@ import java.util.List;
|
||||
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisListCommands;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
/*
|
||||
* Copyright 2017 the original author or 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.springframework.data.redis.connection.lettuce;
|
||||
|
||||
import io.lettuce.core.protocol.RedisCommand;
|
||||
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.data.redis.connection.FutureResult;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
/**
|
||||
* Lettuce specific {@link FutureResult} implementation. <br />
|
||||
*
|
||||
* @author Costin Leau
|
||||
* @author Jennifer Hickey
|
||||
* @author Christoph Strobl
|
||||
* @author Mark Paluch
|
||||
* @since 2.1
|
||||
*/
|
||||
@SuppressWarnings("rawtypes")
|
||||
class LettuceResult<T, S> extends FutureResult<RedisCommand<?, T, ?>> {
|
||||
|
||||
private final boolean convertPipelineAndTxResults;
|
||||
|
||||
<T> LettuceResult(Future<T> resultHolder) {
|
||||
this(resultHolder, false, val -> val);
|
||||
}
|
||||
|
||||
<T> LettuceResult(Future<T> resultHolder, boolean convertPipelineAndTxResults, @Nullable Converter<T, ?> converter) {
|
||||
this(resultHolder, () -> null, convertPipelineAndTxResults, converter);
|
||||
}
|
||||
|
||||
<T> LettuceResult(Future<T> resultHolder, Supplier<S> defaultReturnValue, boolean convertPipelineAndTxResults,
|
||||
@Nullable Converter<T, ?> converter) {
|
||||
|
||||
super((RedisCommand) resultHolder, converter, defaultReturnValue);
|
||||
this.convertPipelineAndTxResults = convertPipelineAndTxResults;
|
||||
}
|
||||
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.redis.connection.FutureResult#get()
|
||||
* @return
|
||||
*/
|
||||
@Nullable
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
public T get() {
|
||||
return (T) getResultHolder().getOutput().get();
|
||||
}
|
||||
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.redis.connection.FutureResult#seeksConversion()
|
||||
* @return
|
||||
*/
|
||||
@Override
|
||||
public boolean seeksConversion() {
|
||||
return convertPipelineAndTxResults && converter != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Lettuce specific {@link FutureResult} implementation of a throw away status result.
|
||||
*/
|
||||
static class LettuceStatusResult extends LettuceResult {
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
LettuceStatusResult(Future resultHolder) {
|
||||
super(resultHolder);
|
||||
setStatus(true);
|
||||
}
|
||||
|
||||
<T> LettuceStatusResult(Future<T> resultHolder, boolean convertPipelineAndTxResults, Converter<T, ?> converter) {
|
||||
super(resultHolder, convertPipelineAndTxResults, converter);
|
||||
setStatus(true);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Lettuce specific {@link FutureResult} implementation of a transaction result.
|
||||
*/
|
||||
static class LettuceTxResult<T> extends FutureResult<Object> {
|
||||
|
||||
private final boolean convertPipelineAndTxResults;
|
||||
|
||||
LettuceTxResult(T resultHolder) {
|
||||
this(resultHolder, false, val -> val);
|
||||
}
|
||||
|
||||
LettuceTxResult(T resultHolder, boolean convertPipelineAndTxResults, Converter<?, ?> converter) {
|
||||
this(resultHolder, () -> null, convertPipelineAndTxResults, converter);
|
||||
}
|
||||
|
||||
LettuceTxResult(T resultHolder, Supplier<Object> defaultReturnValue, boolean convertPipelineAndTxResults,
|
||||
Converter<?, ?> converter) {
|
||||
super(resultHolder, converter, defaultReturnValue);
|
||||
this.convertPipelineAndTxResults = convertPipelineAndTxResults;
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
public Object get() {
|
||||
return getResultHolder();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean seeksConversion() {
|
||||
return convertPipelineAndTxResults && converter != null;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* Lettuce specific {@link FutureResult} implementation of a throw away status result.
|
||||
*/
|
||||
static class LettuceTxStatusResult extends LettuceTxResult {
|
||||
|
||||
LettuceTxStatusResult(Object resultHolder) {
|
||||
super(resultHolder);
|
||||
setStatus(true);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Builder for constructing {@link LettuceResult}.
|
||||
*
|
||||
* @param <T>
|
||||
* @param <S>
|
||||
* @since 2.1
|
||||
*/
|
||||
static class LettuceResultBuilder<T, S> {
|
||||
|
||||
private final Object response;
|
||||
private Converter<T, ?> converter;
|
||||
private boolean convertPipelineAndTxResults = false;
|
||||
private Supplier<?> nullValueDefault = () -> null;
|
||||
|
||||
LettuceResultBuilder(Object response) {
|
||||
|
||||
this.response = response;
|
||||
this.converter = (source) -> source;
|
||||
}
|
||||
|
||||
static <T> LettuceResultBuilder<T, ?> forResponse(Future<T> response) {
|
||||
return new LettuceResultBuilder<>(response);
|
||||
}
|
||||
|
||||
static <T> LettuceResultBuilder<T, ?> forResponse(T response) {
|
||||
return new LettuceResultBuilder<>(response);
|
||||
}
|
||||
|
||||
<S> LettuceResultBuilder<T, S> mappedWith(Converter<T, S> converter) {
|
||||
|
||||
this.converter = converter;
|
||||
return (LettuceResultBuilder<T, S>) this;
|
||||
}
|
||||
|
||||
<S> LettuceResultBuilder<T, S> defaultNullTo(S value) {
|
||||
return (defaultNullTo(() -> value));
|
||||
}
|
||||
|
||||
<S> LettuceResultBuilder<T, S> defaultNullTo(Supplier<S> value) {
|
||||
|
||||
this.nullValueDefault = value;
|
||||
return (LettuceResultBuilder<T, S>) this;
|
||||
}
|
||||
|
||||
LettuceResultBuilder<T, S> convertPipelineAndTxResults(boolean flag) {
|
||||
|
||||
convertPipelineAndTxResults = flag;
|
||||
return this;
|
||||
}
|
||||
|
||||
LettuceResult<T, S> build() {
|
||||
return new LettuceResult((Future<T>) response, nullValueDefault, convertPipelineAndTxResults, converter);
|
||||
}
|
||||
|
||||
LettuceTxResult<T> buildTxResult() {
|
||||
|
||||
return new LettuceTxResult(response, nullValueDefault, convertPipelineAndTxResults, converter);
|
||||
}
|
||||
|
||||
LettuceResult buildStatusResult() {
|
||||
return new LettuceStatusResult((Future<T>) response);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -25,8 +25,7 @@ import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisScriptingCommands;
|
||||
import org.springframework.data.redis.connection.ReturnType;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
|
||||
@@ -27,8 +27,7 @@ import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisNode;
|
||||
import org.springframework.data.redis.connection.RedisServerCommands;
|
||||
import org.springframework.data.redis.connection.convert.Converters;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.core.types.RedisClientInfo;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -29,8 +29,7 @@ import java.util.Set;
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisSetCommands;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.KeyBoundCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -27,8 +27,7 @@ import java.util.concurrent.Future;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisStringCommands;
|
||||
import org.springframework.data.redis.connection.convert.Converters;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.core.types.Expiration;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -131,11 +130,13 @@ class LettuceStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newLettuceStatusResult(getAsyncConnection().set(key, value), Converters.stringToBooleanConverter()));
|
||||
pipeline(
|
||||
connection.newLettuceResult(getAsyncConnection().set(key, value), Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newLettuceTxStatusResult(getConnection().set(key, value), Converters.stringToBooleanConverter()));
|
||||
transaction(
|
||||
connection.newLettuceTxResult(getConnection().set(key, value), Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(getConnection().set(key, value));
|
||||
@@ -158,16 +159,19 @@ class LettuceStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newLettuceStatusResult(
|
||||
getAsyncConnection().set(key, value, LettuceConverters.toSetArgs(expiration, option)), Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newLettuceResult(
|
||||
getAsyncConnection().set(key, value, LettuceConverters.toSetArgs(expiration, option)),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newLettuceTxStatusResult(
|
||||
getConnection().set(key, value, LettuceConverters.toSetArgs(expiration, option)), Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newLettuceTxResult(
|
||||
getConnection().set(key, value, LettuceConverters.toSetArgs(expiration, option)),
|
||||
Converters.stringToBooleanConverter(), () -> false));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(getConnection().set(key, value, LettuceConverters.toSetArgs(expiration, option)));
|
||||
return Converters
|
||||
.stringToBoolean(getConnection().set(key, value, LettuceConverters.toSetArgs(expiration, option)));
|
||||
} catch (Exception ex) {
|
||||
throw convertLettuceAccessException(ex);
|
||||
}
|
||||
@@ -210,11 +214,13 @@ class LettuceStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newLettuceStatusResult(getAsyncConnection().setex(key, seconds, value), Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newLettuceResult(getAsyncConnection().setex(key, seconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newLettuceTxStatusResult(getConnection().setex(key, seconds, value), Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newLettuceTxResult(getConnection().setex(key, seconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(getConnection().setex(key, seconds, value));
|
||||
@@ -235,11 +241,13 @@ class LettuceStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newLettuceStatusResult(getAsyncConnection().psetex(key, milliseconds, value), Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newLettuceResult(getAsyncConnection().psetex(key, milliseconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(connection.newLettuceTxStatusResult(getConnection().psetex(key, milliseconds, value), Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newLettuceTxResult(getConnection().psetex(key, milliseconds, value),
|
||||
Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(getConnection().psetex(key, milliseconds, value));
|
||||
@@ -259,13 +267,11 @@ class LettuceStringCommands implements RedisStringCommands {
|
||||
|
||||
try {
|
||||
if (isPipelined()) {
|
||||
pipeline(connection.newLettuceStatusResult(getAsyncConnection().mset(tuples),
|
||||
Converters.stringToBooleanConverter()));
|
||||
pipeline(connection.newLettuceResult(getAsyncConnection().mset(tuples), Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
if (isQueueing()) {
|
||||
transaction(
|
||||
connection.newLettuceTxStatusResult(getConnection().mset(tuples), Converters.stringToBooleanConverter()));
|
||||
transaction(connection.newLettuceTxResult(getConnection().mset(tuples), Converters.stringToBooleanConverter()));
|
||||
return null;
|
||||
}
|
||||
return Converters.stringToBoolean(getConnection().mset(tuples));
|
||||
|
||||
@@ -29,8 +29,7 @@ import java.util.Set;
|
||||
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.data.redis.connection.RedisZSetCommands;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceConnection.LettuceTxResult;
|
||||
import org.springframework.data.redis.connection.lettuce.LettuceResult.LettuceTxResult;
|
||||
import org.springframework.data.redis.core.Cursor;
|
||||
import org.springframework.data.redis.core.KeyBoundCursor;
|
||||
import org.springframework.data.redis.core.ScanIteration;
|
||||
|
||||
@@ -155,27 +155,31 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test
|
||||
@IfProfileValue(name = "runLongTests", value = "true")
|
||||
public void testExpire() throws Exception {
|
||||
connection.set("exp", "true");
|
||||
|
||||
actual.add(connection.set("exp", "true"));
|
||||
actual.add(connection.expire("exp", 1));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
|
||||
verifyResults(Arrays.asList(true, true));
|
||||
assertTrue(waitFor(new KeyExpired("exp"), 3000l));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "runLongTests", value = "true")
|
||||
public void testExpireAt() throws Exception {
|
||||
connection.set("exp2", "true");
|
||||
|
||||
actual.add(connection.set("exp2", "true"));
|
||||
actual.add(connection.expireAt("exp2", System.currentTimeMillis() / 1000 + 1));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
verifyResults(Arrays.asList(true, true));
|
||||
assertTrue(waitFor(new KeyExpired("exp2"), 3000l));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testPExpire() {
|
||||
connection.set("exp", "true");
|
||||
|
||||
actual.add(connection.set("exp", "true"));
|
||||
actual.add(connection.pExpire("exp", 100));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
verifyResults(Arrays.asList(true, true));
|
||||
assertTrue(waitFor(new KeyExpired("exp"), 1000l));
|
||||
}
|
||||
|
||||
@@ -189,9 +193,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testPExpireAt() {
|
||||
connection.set("exp2", "true");
|
||||
|
||||
actual.add(connection.set("exp2", "true"));
|
||||
actual.add(connection.pExpireAt("exp2", System.currentTimeMillis() + 200));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
verifyResults(Arrays.asList(true, true));
|
||||
assertTrue(waitFor(new KeyExpired("exp2"), 1000l));
|
||||
}
|
||||
|
||||
@@ -397,19 +402,22 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test
|
||||
@IfProfileValue(name = "runLongTests", value = "true")
|
||||
public void testPersist() throws Exception {
|
||||
connection.set("exp3", "true");
|
||||
|
||||
actual.add(connection.set("exp3", "true"));
|
||||
actual.add(connection.expire("exp3", 30));
|
||||
actual.add(connection.persist("exp3"));
|
||||
actual.add(connection.ttl("exp3"));
|
||||
verifyResults(Arrays.asList(new Object[] { true, true, -1L }));
|
||||
verifyResults(Arrays.asList(true, true, true, -1L));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "runLongTests", value = "true")
|
||||
public void testSetEx() throws Exception {
|
||||
connection.setEx("expy", 1l, "yep");
|
||||
|
||||
actual.add(connection.setEx("expy", 1l, "yep"));
|
||||
actual.add(connection.get("expy"));
|
||||
verifyResults(Arrays.asList(new Object[] { "yep" }));
|
||||
|
||||
verifyResults(Arrays.asList(true, "yep"));
|
||||
assertTrue(waitFor(new KeyExpired("expy"), 2500l));
|
||||
}
|
||||
|
||||
@@ -417,10 +425,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@IfProfileValue(name = "runLongTests", value = "true")
|
||||
public void testPsetEx() throws Exception {
|
||||
|
||||
connection.pSetEx("expy", 500L, "yep");
|
||||
actual.add(connection.pSetEx("expy", 500L, "yep"));
|
||||
actual.add(connection.get("expy"));
|
||||
|
||||
verifyResults(Arrays.asList(new Object[] { "yep" }));
|
||||
verifyResults(Arrays.asList(true, "yep"));
|
||||
assertTrue(waitFor(new KeyExpired("expy"), 2500L));
|
||||
}
|
||||
|
||||
@@ -447,11 +455,13 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testSetAndGet() {
|
||||
|
||||
String key = "foo";
|
||||
String value = "blabla";
|
||||
connection.set(key.getBytes(), value.getBytes());
|
||||
|
||||
actual.add(connection.set(key.getBytes(), value.getBytes()));
|
||||
actual.add(connection.get(key));
|
||||
verifyResults(new ArrayList<>(Collections.singletonList(value)));
|
||||
verifyResults(Arrays.asList(true, value));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -485,9 +495,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testBitCountInterval() {
|
||||
connection.set("mykey", "foobar");
|
||||
|
||||
actual.add(connection.set("mykey", "foobar"));
|
||||
actual.add(connection.bitCount("mykey", 1, 1));
|
||||
verifyResults(new ArrayList<>(Collections.singletonList(6l)));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, 6L));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -500,51 +511,57 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testBitOpAnd() {
|
||||
connection.set("key1", "foo");
|
||||
connection.set("key2", "bar");
|
||||
|
||||
actual.add(connection.set("key1", "foo"));
|
||||
actual.add(connection.set("key2", "bar"));
|
||||
actual.add(connection.bitOp(BitOperation.AND, "key3", "key1", "key2"));
|
||||
actual.add(connection.get("key3"));
|
||||
verifyResults(Arrays.asList(new Object[] { 3l, "bab" }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, Boolean.TRUE, 3L, "bab"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testBitOpOr() {
|
||||
connection.set("key1", "foo");
|
||||
connection.set("key2", "ugh");
|
||||
|
||||
actual.add(connection.set("key1", "foo"));
|
||||
actual.add(connection.set("key2", "ugh"));
|
||||
actual.add(connection.bitOp(BitOperation.OR, "key3", "key1", "key2"));
|
||||
actual.add(connection.get("key3"));
|
||||
verifyResults(Arrays.asList(new Object[] { 3l, "woo" }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, Boolean.TRUE, 3l, "woo"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testBitOpXOr() {
|
||||
connection.set("key1", "abcd");
|
||||
connection.set("key2", "efgh");
|
||||
|
||||
actual.add(connection.set("key1", "abcd"));
|
||||
actual.add(connection.set("key2", "efgh"));
|
||||
actual.add(connection.bitOp(BitOperation.XOR, "key3", "key1", "key2"));
|
||||
verifyResults(Arrays.asList(new Object[] { 4l }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, Boolean.TRUE, 4L));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testBitOpNot() {
|
||||
connection.set("key1", "abcd");
|
||||
|
||||
actual.add(connection.set("key1", "abcd"));
|
||||
actual.add(connection.bitOp(BitOperation.NOT, "key3", "key1"));
|
||||
verifyResults(Arrays.asList(new Object[] { 4l }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, 4L));
|
||||
}
|
||||
|
||||
@Test(expected = UnsupportedOperationException.class)
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testBitOpNotMultipleSources() {
|
||||
connection.set("key1", "abcd");
|
||||
connection.set("key2", "efgh");
|
||||
|
||||
actual.add(connection.set("key1", "abcd"));
|
||||
actual.add(connection.set("key2", "efgh"));
|
||||
actual.add(connection.bitOp(BitOperation.NOT, "key3", "key1", "key2"));
|
||||
getResults();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testInfo() throws Exception {
|
||||
|
||||
actual.add(connection.info());
|
||||
List<Object> results = getResults();
|
||||
Properties info = (Properties) results.get(0);
|
||||
@@ -629,10 +646,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testAppend() {
|
||||
connection.set("a", "b");
|
||||
actual.add(connection.set("a", "b"));
|
||||
actual.add(connection.append("a", "c"));
|
||||
actual.add(connection.get("a"));
|
||||
verifyResults(Arrays.asList(new Object[] { 2l, "bc" }));
|
||||
verifyResults(Arrays.asList(new Object[] { Boolean.TRUE, 2l, "bc" }));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -731,13 +748,16 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testExecute() {
|
||||
connection.set("foo", "bar");
|
||||
|
||||
actual.add(connection.set("foo", "bar"));
|
||||
actual.add(connection.execute("GET", "foo"));
|
||||
assertEquals("bar", stringSerializer.deserialize((byte[]) getResults().get(0)));
|
||||
|
||||
assertEquals("bar", stringSerializer.deserialize((byte[]) getResults().get(1)));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testExecuteNoArgs() {
|
||||
|
||||
actual.add(connection.execute("PING"));
|
||||
List<Object> results = getResults();
|
||||
assertEquals("PONG", stringSerializer.deserialize((byte[]) results.get(0)));
|
||||
@@ -746,13 +766,14 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
public void testMultiExec() throws Exception {
|
||||
|
||||
connection.multi();
|
||||
connection.set("key", "value");
|
||||
connection.get("key");
|
||||
actual.add(connection.exec());
|
||||
List<Object> results = getResults();
|
||||
List<Object> execResults = (List<Object>) results.get(0);
|
||||
assertEquals(Arrays.asList(new Object[] { "value" }), execResults);
|
||||
assertEquals(Arrays.asList(true, "value"), execResults);
|
||||
assertEquals("value", connection.get("key"));
|
||||
}
|
||||
|
||||
@@ -796,7 +817,9 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testWatch() throws Exception {
|
||||
connection.set("testitnow", "willdo");
|
||||
|
||||
actual.add(connection.set("testitnow", "willdo"));
|
||||
|
||||
connection.watch("testitnow".getBytes());
|
||||
// Give some time for watch to be asynch executed in extending tests
|
||||
Thread.sleep(500);
|
||||
@@ -809,29 +832,33 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
actual.add(connection.get("testitnow"));
|
||||
|
||||
if (connectionFactory instanceof JedisConnectionFactory) {
|
||||
verifyResults(Arrays.asList(new Object[] { Collections.emptyList(), "something" }));
|
||||
verifyResults(Arrays.asList(new Object[] { true, Collections.emptyList(), "something" }));
|
||||
} else {
|
||||
verifyResults(Arrays.asList(new Object[] { null, "something" }));
|
||||
verifyResults(Arrays.asList(new Object[] { true, null, "something" }));
|
||||
}
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
public void testUnwatch() throws Exception {
|
||||
connection.set("testitnow", "willdo");
|
||||
|
||||
actual.add(connection.set("testitnow", "willdo"));
|
||||
|
||||
connection.watch("testitnow".getBytes());
|
||||
connection.unwatch();
|
||||
|
||||
connection.multi();
|
||||
|
||||
// Give some time for unwatch to be asynch executed
|
||||
Thread.sleep(100);
|
||||
DefaultStringRedisConnection conn2 = new DefaultStringRedisConnection(connectionFactory.getConnection());
|
||||
conn2.set("testitnow", "something");
|
||||
|
||||
connection.set("testitnow", "somethingelse");
|
||||
connection.get("testitnow");
|
||||
actual.add(connection.exec());
|
||||
List<Object> results = getResults();
|
||||
List<Object> execResults = (List<Object>) results.get(0);
|
||||
assertEquals(Arrays.asList(new Object[] { "somethingelse" }), execResults);
|
||||
|
||||
verifyResults(Arrays.asList(true, Arrays.asList(true, "somethingelse")));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -874,9 +901,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testDbSize() {
|
||||
connection.set("dbparam", "foo");
|
||||
|
||||
actual.add(connection.set("dbparam", "foo"));
|
||||
actual.add(connection.dbSize());
|
||||
assertTrue((Long) getResults().get(0) > 0);
|
||||
assertTrue((Long) getResults().get(1) > 0);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -902,22 +930,23 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testExists() {
|
||||
connection.set("existent", "true");
|
||||
|
||||
actual.add(connection.set("existent", "true"));
|
||||
actual.add(connection.exists("existent"));
|
||||
actual.add(connection.exists("nonexistent"));
|
||||
verifyResults(Arrays.asList(new Object[] { true, false }));
|
||||
verifyResults(Arrays.asList(true, true, false));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-529
|
||||
public void testExistsWithMultipleKeys() {
|
||||
|
||||
connection.set("exist-1", "true");
|
||||
connection.set("exist-2", "true");
|
||||
connection.set("exist-3", "true");
|
||||
actual.add(connection.set("exist-1", "true"));
|
||||
actual.add(connection.set("exist-2", "true"));
|
||||
actual.add(connection.set("exist-3", "true"));
|
||||
|
||||
actual.add(connection.exists("exist-1", "exist-2", "exist-3", "nonexistent"));
|
||||
|
||||
verifyResults(Arrays.asList(new Object[] { 3L }));
|
||||
verifyResults(Arrays.asList(new Object[] { true, true, true, 3L }));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-529
|
||||
@@ -931,100 +960,105 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test // DATAREDIS-529
|
||||
public void testExistsSameKeyMultipleTimes() {
|
||||
|
||||
connection.set("existent", "true");
|
||||
actual.add(connection.set("existent", "true"));
|
||||
|
||||
actual.add(connection.exists("existent", "existent"));
|
||||
|
||||
verifyResults(Arrays.asList(new Object[] { 2L }));
|
||||
verifyResults(Arrays.asList(true, 2L));
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
public void testKeys() throws Exception {
|
||||
connection.set("keytest", "true");
|
||||
|
||||
actual.add(connection.set("keytest", "true"));
|
||||
actual.add(connection.keys("key*"));
|
||||
assertTrue(((Collection<String>) getResults().get(0)).contains("keytest"));
|
||||
assertTrue(((Collection<String>) getResults().get(1)).contains("keytest"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRandomKey() {
|
||||
connection.set("some", "thing");
|
||||
|
||||
actual.add(connection.set("some", "thing"));
|
||||
actual.add(connection.randomKey());
|
||||
List<Object> results = getResults();
|
||||
assertNotNull(results.get(0));
|
||||
assertNotNull(results.get(1));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRename() {
|
||||
connection.set("renametest", "testit");
|
||||
|
||||
actual.add(connection.set("renametest", "testit"));
|
||||
connection.rename("renametest", "newrenametest");
|
||||
actual.add(connection.get("newrenametest"));
|
||||
actual.add(connection.exists("renametest"));
|
||||
verifyResults(Arrays.asList(new Object[] { "testit", false }));
|
||||
verifyResults(Arrays.asList(true, "testit", false));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRenameNx() {
|
||||
connection.set("nxtest", "testit");
|
||||
|
||||
actual.add(connection.set("nxtest", "testit"));
|
||||
actual.add(connection.renameNX("nxtest", "newnxtest"));
|
||||
actual.add(connection.get("newnxtest"));
|
||||
actual.add(connection.exists("nxtest"));
|
||||
verifyResults(Arrays.asList(new Object[] { true, "testit", false }));
|
||||
verifyResults(Arrays.asList(true, true, "testit", false));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testTtl() {
|
||||
connection.set("whatup", "yo");
|
||||
actual.add(connection.set("whatup", "yo"));
|
||||
actual.add(connection.ttl("whatup"));
|
||||
verifyResults(Arrays.asList(new Object[] { -1L }));
|
||||
verifyResults(Arrays.asList(true, -1L));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-526
|
||||
public void testTtlWithTimeUnit() {
|
||||
|
||||
connection.set("whatup", "yo");
|
||||
actual.add(connection.set("whatup", "yo"));
|
||||
actual.add(connection.expire("whatup", 10));
|
||||
actual.add(connection.ttl("whatup", TimeUnit.MILLISECONDS));
|
||||
|
||||
List<Object> results = getResults();
|
||||
|
||||
assertTrue((Long) results.get(1) > TimeUnit.SECONDS.toMillis(5));
|
||||
assertTrue((Long) results.get(1) <= TimeUnit.SECONDS.toMillis(10));
|
||||
assertTrue((Long) results.get(2) > TimeUnit.SECONDS.toMillis(5));
|
||||
assertTrue((Long) results.get(2) <= TimeUnit.SECONDS.toMillis(10));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testPTtlNoExpire() {
|
||||
connection.set("whatup", "yo");
|
||||
|
||||
actual.add(connection.set("whatup", "yo"));
|
||||
actual.add(connection.pTtl("whatup"));
|
||||
verifyResults(Arrays.asList(new Object[] { -1L }));
|
||||
verifyResults(Arrays.asList(true, -1L));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testPTtl() {
|
||||
|
||||
connection.set("whatup", "yo");
|
||||
actual.add(connection.set("whatup", "yo"));
|
||||
actual.add(connection.pExpire("whatup", TimeUnit.SECONDS.toMillis(10)));
|
||||
actual.add(connection.pTtl("whatup"));
|
||||
|
||||
List<Object> results = getResults();
|
||||
|
||||
assertTrue((Long) results.get(1) > -1);
|
||||
assertTrue((Long) results.get(2) > -1);
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-526
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testPTtlWithTimeUnit() {
|
||||
|
||||
connection.set("whatup", "yo");
|
||||
actual.add(connection.set("whatup", "yo"));
|
||||
actual.add(connection.pExpire("whatup", TimeUnit.MINUTES.toMillis(10)));
|
||||
actual.add(connection.pTtl("whatup", TimeUnit.SECONDS));
|
||||
|
||||
List<Object> results = getResults();
|
||||
|
||||
assertTrue((Long) results.get(1) > TimeUnit.MINUTES.toSeconds(9));
|
||||
assertTrue((Long) results.get(1) <= TimeUnit.MINUTES.toSeconds(10));
|
||||
assertTrue((Long) results.get(2) > TimeUnit.MINUTES.toSeconds(9));
|
||||
assertTrue((Long) results.get(2) <= TimeUnit.MINUTES.toSeconds(10));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1062,49 +1096,54 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
@Test(expected = RedisSystemException.class)
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testRestoreExistingKey() {
|
||||
connection.set("testing", "12");
|
||||
|
||||
actual.add(connection.set("testing", "12"));
|
||||
actual.add(connection.dump("testing".getBytes()));
|
||||
List<Object> results = getResults();
|
||||
initConnection();
|
||||
connection.restore("testing".getBytes(), 0, (byte[]) results.get(0));
|
||||
connection.restore("testing".getBytes(), 0, (byte[]) results.get(1));
|
||||
getResults();
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testRestoreTtl() {
|
||||
connection.set("testing", "12");
|
||||
|
||||
actual.add(connection.set("testing", "12"));
|
||||
actual.add(connection.dump("testing".getBytes()));
|
||||
|
||||
List<Object> results = getResults();
|
||||
initConnection();
|
||||
actual.add(connection.del("testing"));
|
||||
actual.add(connection.get("testing"));
|
||||
connection.restore("testing".getBytes(), 100l, (byte[]) results.get(0));
|
||||
verifyResults(Arrays.asList(new Object[] { 1l, null }));
|
||||
connection.restore("testing".getBytes(), 100l, (byte[]) results.get(1));
|
||||
verifyResults(Arrays.asList(1l, null));
|
||||
assertTrue(waitFor(new KeyExpired("testing"), 400l));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDel() {
|
||||
connection.set("testing", "123");
|
||||
|
||||
actual.add(connection.set("testing", "123"));
|
||||
actual.add(connection.del("testing"));
|
||||
actual.add(connection.exists("testing"));
|
||||
verifyResults(Arrays.asList(new Object[] { 1l, false }));
|
||||
verifyResults(Arrays.asList(true, 1L, false));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testType() {
|
||||
connection.set("something", "yo");
|
||||
|
||||
actual.add(connection.set("something", "yo"));
|
||||
actual.add(connection.type("something"));
|
||||
verifyResults(Arrays.asList(new Object[] { DataType.STRING }));
|
||||
verifyResults(Arrays.asList(true, DataType.STRING));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetSet() {
|
||||
connection.set("testGS", "1");
|
||||
actual.add(connection.set("testGS", "1"));
|
||||
actual.add(connection.getSet("testGS", "2"));
|
||||
actual.add(connection.get("testGS"));
|
||||
verifyResults(Arrays.asList(new Object[] { "1", "2" }));
|
||||
verifyResults(Arrays.asList(true, "1", "2"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1112,9 +1151,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
Map<String, String> vals = new HashMap<>();
|
||||
vals.put("color", "orange");
|
||||
vals.put("size", "1");
|
||||
connection.mSetString(vals);
|
||||
|
||||
actual.add(connection.mSetString(vals));
|
||||
actual.add(connection.mGet("color", "size"));
|
||||
verifyResults(Arrays.asList(new Object[] { Arrays.asList(new String[] { "orange", "1" }) }));
|
||||
verifyResults(Arrays.asList(true, Arrays.asList(new String[] { "orange", "1" })));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1129,13 +1169,13 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testMSetNxFailure() {
|
||||
connection.set("height", "2");
|
||||
actual.add(connection.set("height", "2"));
|
||||
Map<String, String> vals = new HashMap<>();
|
||||
vals.put("height", "5");
|
||||
vals.put("width", "1");
|
||||
actual.add(connection.mSetNXString(vals));
|
||||
actual.add(connection.mGet("height", "width"));
|
||||
verifyResults(Arrays.asList(new Object[] { false, Arrays.asList(new String[] { "2", null }) }));
|
||||
verifyResults(Arrays.asList(true, false, Arrays.asList(new String[] { "2", null })));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1149,43 +1189,49 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testGetRangeSetRange() {
|
||||
connection.set("rangekey", "supercalifrag");
|
||||
|
||||
actual.add(connection.set("rangekey", "supercalifrag"));
|
||||
actual.add(connection.getRange("rangekey", 0l, 2l));
|
||||
connection.setRange("rangekey", "ck", 2);
|
||||
actual.add(connection.get("rangekey"));
|
||||
verifyResults(Arrays.asList(new Object[] { "sup", "suckrcalifrag" }));
|
||||
verifyResults(Arrays.asList(true, "sup", "suckrcalifrag"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDecrByIncrBy() {
|
||||
connection.set("tdb", "4");
|
||||
|
||||
actual.add(connection.set("tdb", "4"));
|
||||
actual.add(connection.decrBy("tdb", 3l));
|
||||
actual.add(connection.incrBy("tdb", 7l));
|
||||
verifyResults(Arrays.asList(new Object[] { 1l, 8l }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, 1L, 8L));
|
||||
}
|
||||
|
||||
@Test
|
||||
@IfProfileValue(name = "redisVersion", value = "2.6+")
|
||||
public void testIncrByDouble() {
|
||||
connection.set("tdb", "4.5");
|
||||
|
||||
actual.add(connection.set("tdb", "4.5"));
|
||||
actual.add(connection.incrBy("tdb", 7.2));
|
||||
actual.add(connection.get("tdb"));
|
||||
verifyResults(Arrays.asList(new Object[] { 11.7d, "11.7" }));
|
||||
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, 11.7d, "11.7"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testIncrDecrByLong() {
|
||||
|
||||
String key = "test.count";
|
||||
long largeNumber = 0x123456789L; // > 32bits
|
||||
connection.set(key, "0");
|
||||
actual.add(connection.set(key, "0"));
|
||||
actual.add(connection.incrBy(key, largeNumber));
|
||||
actual.add(connection.decrBy(key, largeNumber));
|
||||
actual.add(connection.decrBy(key, 2 * largeNumber));
|
||||
verifyResults(Arrays.asList(new Object[] { largeNumber, 0l, -2 * largeNumber }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, largeNumber, 0l, -2 * largeNumber));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testHashIncrDecrByLong() {
|
||||
|
||||
String key = "test.hcount";
|
||||
String hkey = "hashkey";
|
||||
|
||||
@@ -1200,19 +1246,21 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testIncDecr() {
|
||||
connection.set("incrtest", "0");
|
||||
|
||||
actual.add(connection.set("incrtest", "0"));
|
||||
actual.add(connection.incr("incrtest"));
|
||||
actual.add(connection.get("incrtest"));
|
||||
actual.add(connection.decr("incrtest"));
|
||||
actual.add(connection.get("incrtest"));
|
||||
verifyResults(Arrays.asList(new Object[] { 1l, "1", 0l, "0" }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, 1L, "1", 0L, "0"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testStrLen() {
|
||||
connection.set("strlentest", "cat");
|
||||
|
||||
actual.add(connection.set("strlentest", "cat"));
|
||||
actual.add(connection.strLen("strlentest"));
|
||||
verifyResults(Arrays.asList(new Object[] { 3l }));
|
||||
verifyResults(Arrays.asList(Boolean.TRUE, 3L));
|
||||
}
|
||||
|
||||
// List operations
|
||||
@@ -1943,9 +1991,11 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void testMove() {
|
||||
connection.set("foo", "bar");
|
||||
|
||||
actual.add(connection.set("foo", "bar"));
|
||||
actual.add(connection.move("foo", 1));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
|
||||
verifyResults(Arrays.asList(true, true));
|
||||
connection.select(1);
|
||||
try {
|
||||
assertEquals("bar", connection.get("foo"));
|
||||
@@ -2250,14 +2300,15 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithExpirationAndUpsertOpionShouldSetTtlWhenKeyDoesNotExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "foo", Expiration.milliseconds(500), SetOption.upsert());
|
||||
actual.add(connection.set(key, "foo", Expiration.milliseconds(500), SetOption.upsert()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(500d, 499d)));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(2)).doubleValue(), is(closeTo(500d, 499d)));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2265,8 +2316,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithExpirationAndUpsertOpionShouldSetTtlWhenKeyDoesExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "spring");
|
||||
connection.set(key, "data", Expiration.milliseconds(500), SetOption.upsert());
|
||||
actual.add(connection.set(key, "spring"));
|
||||
actual.add(connection.set(key, "data", Expiration.milliseconds(500), SetOption.upsert()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2274,8 +2325,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(500d, 499d)));
|
||||
assertThat((result.get(2)), is(equalTo("data")));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(result.get(2), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(3)).doubleValue(), is(closeTo(500d, 499d)));
|
||||
assertThat((result.get(4)), is(equalTo("data")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2283,8 +2336,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithExpirationAndAbsentOptionShouldSetTtlWhenKeyDoesExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "spring");
|
||||
connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifAbsent());
|
||||
actual.add(connection.set(key, "spring"));
|
||||
actual.add(connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifAbsent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2292,8 +2345,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(2)), is(equalTo("spring")));
|
||||
assertThat(result.get(1), is(Boolean.FALSE));
|
||||
assertThat(result.get(2), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(3)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(4)), is(equalTo("spring")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2301,7 +2356,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithExpirationAndAbsentOptionShouldSetTtlWhenKeyDoesNotExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifAbsent());
|
||||
actual.add(connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifAbsent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2309,8 +2364,9 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(500d, 499d)));
|
||||
assertThat((result.get(2)), is(equalTo("data")));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(2)).doubleValue(), is(closeTo(500d, 499d)));
|
||||
assertThat((result.get(3)), is(equalTo("data")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2318,8 +2374,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithExpirationAndPresentOptionShouldSetTtlWhenKeyDoesExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "spring");
|
||||
connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifPresent());
|
||||
actual.add(connection.set(key, "spring"));
|
||||
actual.add(connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifPresent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2327,8 +2383,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(500, 499)));
|
||||
assertThat((result.get(2)), is(equalTo("data")));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(result.get(2), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(3)).doubleValue(), is(closeTo(500, 499)));
|
||||
assertThat((result.get(4)), is(equalTo("data")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2336,7 +2394,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithExpirationAndPresentOptionShouldSetTtlWhenKeyDoesNotExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifPresent());
|
||||
actual.add(connection.set(key, "data", Expiration.milliseconds(500), SetOption.ifPresent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2344,7 +2402,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.FALSE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-2, 0)));
|
||||
assertThat(result.get(1), is(Boolean.FALSE));
|
||||
assertThat(((Long) result.get(2)).doubleValue(), is(closeTo(-2, 0)));
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class) // DATAREDIS-316, DATAREDIS-692
|
||||
@@ -2360,14 +2419,15 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithoutExpirationAndUpsertOpionShouldSetTtlWhenKeyDoesNotExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "foo", Expiration.persistent(), SetOption.upsert());
|
||||
actual.add(connection.set(key, "foo", Expiration.persistent(), SetOption.upsert()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(2)).doubleValue(), is(closeTo(-1, 0)));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2375,17 +2435,20 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithoutExpirationAndUpsertOpionShouldSetTtlWhenKeyDoesExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "spring");
|
||||
connection.set(key, "data", Expiration.persistent(), SetOption.upsert());
|
||||
actual.add(connection.set(key, "spring"));
|
||||
actual.add(connection.set(key, "data", Expiration.persistent(), SetOption.upsert()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
actual.add(connection.get(key));
|
||||
|
||||
List<Object> result = getResults();
|
||||
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(2)), is(equalTo("data")));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(result.get(2), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(3)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(4)), is(equalTo("data")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2393,8 +2456,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithoutExpirationAndAbsentOptionShouldSetTtlWhenKeyDoesExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "spring");
|
||||
connection.set(key, "data", Expiration.persistent(), SetOption.ifAbsent());
|
||||
actual.add(connection.set(key, "spring"));
|
||||
actual.add(connection.set(key, "data", Expiration.persistent(), SetOption.ifAbsent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2402,8 +2465,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(2)), is(equalTo("spring")));
|
||||
assertThat(result.get(1), is(Boolean.FALSE));
|
||||
assertThat(result.get(2), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(3)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(4)), is(equalTo("spring")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2411,7 +2476,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithoutExpirationAndAbsentOptionShouldSetTtlWhenKeyDoesNotExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "data", Expiration.persistent(), SetOption.ifAbsent());
|
||||
actual.add(connection.set(key, "data", Expiration.persistent(), SetOption.ifAbsent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2419,8 +2484,9 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(2)), is(equalTo("data")));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(2)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(3)), is(equalTo("data")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2428,8 +2494,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithoutExpirationAndPresentOptionShouldSetTtlWhenKeyDoesExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "spring");
|
||||
connection.set(key, "data", Expiration.persistent(), SetOption.ifPresent());
|
||||
actual.add(connection.set(key, "spring"));
|
||||
actual.add(connection.set(key, "data", Expiration.persistent(), SetOption.ifPresent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2437,8 +2503,10 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(2)), is(equalTo("data")));
|
||||
assertThat(result.get(1), is(Boolean.TRUE));
|
||||
assertThat(result.get(2), is(Boolean.TRUE));
|
||||
assertThat(((Long) result.get(3)).doubleValue(), is(closeTo(-1, 0)));
|
||||
assertThat((result.get(4)), is(equalTo("data")));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-316
|
||||
@@ -2446,7 +2514,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
public void setWithoutExpirationAndPresentOptionShouldSetTtlWhenKeyDoesNotExist() {
|
||||
|
||||
String key = "exp-" + UUID.randomUUID();
|
||||
connection.set(key, "data", Expiration.persistent(), SetOption.ifPresent());
|
||||
actual.add(connection.set(key, "data", Expiration.persistent(), SetOption.ifPresent()));
|
||||
|
||||
actual.add(connection.exists(key));
|
||||
actual.add(connection.pTtl(key));
|
||||
@@ -2454,7 +2522,8 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
List<Object> result = getResults();
|
||||
assertThat(result.get(0), is(Boolean.FALSE));
|
||||
assertThat(((Long) result.get(1)).doubleValue(), is(closeTo(-2, 0)));
|
||||
assertThat(result.get(1), is(Boolean.FALSE));
|
||||
assertThat(((Long) result.get(2)).doubleValue(), is(closeTo(-2, 0)));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-438
|
||||
@@ -2685,8 +2754,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
actual.add(connection.hSet("hash-hstrlen", "key-2", "value-2"));
|
||||
actual.add(connection.hStrLen("hash-hstrlen", "key-2"));
|
||||
|
||||
verifyResults(
|
||||
Arrays.asList(new Object[] { Boolean.TRUE, Boolean.TRUE, Long.valueOf("value-2".length()) }));
|
||||
verifyResults(Arrays.asList(new Object[] { Boolean.TRUE, Boolean.TRUE, Long.valueOf("value-2".length()) }));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-698
|
||||
@@ -2695,8 +2763,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
actual.add(connection.hSet("hash-hstrlen", "key-1", "value-1"));
|
||||
actual.add(connection.hStrLen("hash-hstrlen", "key-2"));
|
||||
|
||||
verifyResults(
|
||||
Arrays.asList(new Object[] { Boolean.TRUE, 0L }));
|
||||
verifyResults(Arrays.asList(new Object[] { Boolean.TRUE, 0L }));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-698
|
||||
@@ -2704,8 +2771,7 @@ public abstract class AbstractConnectionIntegrationTests {
|
||||
|
||||
actual.add(connection.hStrLen("hash-no-exist", "key-2"));
|
||||
|
||||
verifyResults(
|
||||
Arrays.asList(new Object[] { 0L }));
|
||||
verifyResults(Arrays.asList(new Object[] { 0L }));
|
||||
}
|
||||
|
||||
protected void verifyResults(List<Object> expected) {
|
||||
|
||||
@@ -102,7 +102,7 @@ public class JedisConnectionPipelineIntegrationTests extends AbstractConnectionP
|
||||
actual.add(connection.exec());
|
||||
List<Object> results = getResults();
|
||||
List<Object> execResults = (List<Object>) results.get(0);
|
||||
assertEquals(Arrays.asList(new Object[] { "somethingelse" }), execResults);
|
||||
assertEquals(Arrays.asList(new Object[] { true, "somethingelse" }), execResults);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -83,9 +83,10 @@ public class LettuceConnectionPipelineIntegrationTests extends AbstractConnectio
|
||||
|
||||
@Test
|
||||
public void testMove() {
|
||||
connection.set("foo", "bar");
|
||||
|
||||
actual.add(connection.set("foo", "bar"));
|
||||
actual.add(connection.move("foo", 1));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
verifyResults(Arrays.asList(new Object[] { true, true }));
|
||||
// Lettuce does not support select when using shared conn, use a new conn factory
|
||||
LettuceConnectionFactory factory2 = new LettuceConnectionFactory();
|
||||
factory2.setClientResources(LettuceTestClientResources.getSharedClientResources());
|
||||
|
||||
@@ -41,9 +41,11 @@ public class LettuceConnectionTransactionIntegrationTests extends AbstractConnec
|
||||
|
||||
@Test
|
||||
public void testMove() {
|
||||
connection.set("foo", "bar");
|
||||
|
||||
actual.add(connection.set("foo", "bar"));
|
||||
actual.add(connection.move("foo", 1));
|
||||
verifyResults(Arrays.asList(new Object[] { true }));
|
||||
verifyResults(Arrays.asList(true, true ));
|
||||
|
||||
// Lettuce does not support select when using shared conn, use a new conn factory
|
||||
LettuceConnectionFactory factory2 = new LettuceConnectionFactory();
|
||||
factory2.setClientResources(LettuceTestClientResources.getSharedClientResources());
|
||||
|
||||
@@ -184,7 +184,7 @@ public class RedisTemplateTests<K, V> {
|
||||
Set<V> set = new HashSet<>(Collections.singletonList(setValue));
|
||||
Set<TypedTuple<V>> tupleSet = new LinkedHashSet<>(
|
||||
Collections.singletonList(new DefaultTypedTuple<>(zsetValue, 1d)));
|
||||
assertThat(results, isEqual(Arrays.asList(new Object[] { value1, 1l, list, 1l, set, true, tupleSet })));
|
||||
assertThat(results, isEqual(Arrays.asList(new Object[] { true, value1, 1l, list, 1l, set, true, tupleSet })));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -235,7 +235,7 @@ public class RedisTemplateTests<K, V> {
|
||||
Map<Long, Long> map = new LinkedHashMap<>();
|
||||
map.put(10l, 11l);
|
||||
assertThat(results,
|
||||
isEqual(Arrays.asList(new Object[] { 5l, 1L, 1l, list, 1l, longSet, true, tupleSet, zSet, true, map })));
|
||||
isEqual(Arrays.asList(new Object[] { true, 5l, 1L, 1l, list, 1l, longSet, true, tupleSet, zSet, true, map })));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -282,7 +282,7 @@ public class RedisTemplateTests<K, V> {
|
||||
return null;
|
||||
});
|
||||
assertThat(results,
|
||||
isEqual(Arrays.asList(new Object[] { value1, 1l, 2l, Arrays.asList(new Object[] { listValue, listValue2 }) })));
|
||||
isEqual(Arrays.asList(new Object[] { true, value1, 1l, 2l, Arrays.asList(new Object[] { listValue, listValue2 }) })));
|
||||
}
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
@@ -301,7 +301,7 @@ public class RedisTemplateTests<K, V> {
|
||||
return null;
|
||||
}, new GenericToStringSerializer<>(Long.class));
|
||||
|
||||
assertEquals(Arrays.asList(new Object[] { 5l, 1l, 2l, Arrays.asList(new Long[] { 10l, 11l }) }), results);
|
||||
assertEquals(Arrays.asList(new Object[] { true, 5l, 1l, 2l, Arrays.asList(new Long[] { 10l, 11l }) }), results);
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-500
|
||||
@@ -349,7 +349,7 @@ public class RedisTemplateTests<K, V> {
|
||||
});
|
||||
// Should contain the List of deserialized exec results and the result of the last call to get()
|
||||
assertThat(pipelinedResults,
|
||||
isEqual(Arrays.asList(new Object[] { Arrays.asList(new Object[] { 1l, value1, 0l }), value1 })));
|
||||
isEqual(Arrays.asList(new Object[] { Arrays.asList(new Object[] { 1l, value1, 0l }), true, value1 })));
|
||||
}
|
||||
|
||||
@SuppressWarnings({ "rawtypes", "unchecked" })
|
||||
@@ -369,7 +369,7 @@ public class RedisTemplateTests<K, V> {
|
||||
}
|
||||
}, new GenericToStringSerializer<>(Long.class));
|
||||
// Should contain the List of deserialized exec results and the result of the last call to get()
|
||||
assertEquals(Arrays.asList(new Object[] { Arrays.asList(new Object[] { 1l, 5l, 0l }), 2l }), pipelinedResults);
|
||||
assertEquals(Arrays.asList(new Object[] { Arrays.asList(new Object[] { 1l, 5l, 0l }), true, 2l }), pipelinedResults);
|
||||
}
|
||||
|
||||
@Test(expected = InvalidDataAccessApiUsageException.class)
|
||||
@@ -552,9 +552,9 @@ public class RedisTemplateTests<K, V> {
|
||||
}
|
||||
});
|
||||
|
||||
assertThat(result, hasSize(2));
|
||||
assertThat(((Long) result.get(1)), greaterThanOrEqualTo(23L));
|
||||
assertThat(((Long) result.get(1)), lessThan(25L));
|
||||
assertThat(result, hasSize(3));
|
||||
assertThat(((Long) result.get(2)), greaterThanOrEqualTo(23L));
|
||||
assertThat(((Long) result.get(2)), lessThan(25L));
|
||||
}
|
||||
|
||||
@Test // DATAREDIS-526
|
||||
@@ -578,9 +578,9 @@ public class RedisTemplateTests<K, V> {
|
||||
}
|
||||
});
|
||||
|
||||
assertThat(result, hasSize(2));
|
||||
assertThat(((Long) result.get(1)), greaterThanOrEqualTo(23L));
|
||||
assertThat(((Long) result.get(1)), lessThan(25L));
|
||||
assertThat(result, hasSize(3));
|
||||
assertThat(((Long) result.get(2)), greaterThanOrEqualTo(23L));
|
||||
assertThat(((Long) result.get(2)), lessThan(25L));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -741,7 +741,7 @@ public class RedisTemplateTests<K, V> {
|
||||
}
|
||||
});
|
||||
|
||||
assertTrue(results.isEmpty());
|
||||
assertTrue(results.size() == 1);
|
||||
assertThat(redisTemplate.opsForValue().get(key1), isEqual(value3));
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user