pull/5427/head
Manuel Polo 1 year ago
commit eb5eefd034
No known key found for this signature in database
GPG Key ID: 225EFA9B34C561E3

@ -175,14 +175,14 @@
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>testcontainers-bom</artifactId>
<version>1.18.3</version>
<version>1.19.1</version>
<type>pom</type>
<scope>import</scope>
</dependency>
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-bom</artifactId>
<version>4.1.100.Final</version>
<version>4.1.101.Final</version>
<type>pom</type>
<scope>import</scope>
</dependency>

@ -35,7 +35,7 @@
<plugin>
<groupId>com.mycila</groupId>
<artifactId>license-maven-plugin</artifactId>
<version>3.0</version>
<version>4.3</version>
<configuration>
<basedir>${basedir}</basedir>
<header>${basedir}/../header.txt</header>

@ -85,7 +85,7 @@
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-core</artifactId>
<version>3.5.3</version>
<version>3.5.11</version>
</dependency>
<dependency>
<groupId>org.reactivestreams</groupId>
@ -489,7 +489,7 @@
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-checkstyle-plugin</artifactId>
<version>3.2.1</version>
<version>3.3.1</version>
<executions>
<execution>
<phase>verify</phase>
@ -500,7 +500,6 @@
</executions>
<configuration>
<consoleOutput>true</consoleOutput>
<enableRSS>false</enableRSS>
<configLocation>/checkstyle.xml</configLocation>
<propertyExpansion>checkstyle.config.path=${basedir}</propertyExpansion>
</configuration>
@ -508,7 +507,7 @@
<dependency>
<groupId>com.puppycrawl.tools</groupId>
<artifactId>checkstyle</artifactId>
<version>10.8.0</version>
<version>10.12.4</version>
</dependency>
</dependencies>
</plugin>
@ -599,7 +598,7 @@
<plugin>
<groupId>com.mycila</groupId>
<artifactId>license-maven-plugin</artifactId>
<version>2.11</version>
<version>4.3</version>
<configuration>
<basedir>${basedir}</basedir>
<header>${basedir}/../header.txt</header>

@ -0,0 +1,897 @@
/**
* Copyright (c) 2013-2022 Nikita Koksharov
*
* 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.redisson;
import org.redisson.api.*;
import org.redisson.api.listener.*;
import org.redisson.api.mapreduce.RCollectionMapReduce;
import org.redisson.client.RedisClient;
import org.redisson.client.RedisException;
import org.redisson.client.codec.Codec;
import org.redisson.client.codec.StringCodec;
import org.redisson.client.protocol.RedisCommand;
import org.redisson.client.protocol.RedisCommands;
import org.redisson.client.protocol.convertor.BooleanNumberReplayConvertor;
import org.redisson.client.protocol.convertor.Convertor;
import org.redisson.client.protocol.convertor.IntegerReplayConvertor;
import org.redisson.command.CommandAsyncExecutor;
import org.redisson.iterator.RedissonBaseIterator;
import org.redisson.iterator.RedissonListIterator;
import org.redisson.mapreduce.RedissonCollectionMapReduce;
import org.redisson.misc.CompletableFutureWrapper;
import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.CompletionStage;
import java.util.function.Predicate;
import static org.redisson.client.protocol.RedisCommands.*;
/**
* Base list implementation
*
* @author Nikita Koksharov
*
* @param <V> the type of elements held in this collection
*/
public class BaseRedissonList<V> extends RedissonExpirable {
private RedissonClient redisson;
BaseRedissonList(CommandAsyncExecutor commandExecutor, String name, RedissonClient redisson) {
super(commandExecutor, name);
this.redisson = redisson;
}
BaseRedissonList(Codec codec, CommandAsyncExecutor commandExecutor, String name, RedissonClient redisson) {
super(codec, commandExecutor, name);
this.redisson = redisson;
}
public <KOut, VOut> RCollectionMapReduce<V, KOut, VOut> mapReduce() {
return new RedissonCollectionMapReduce<V, KOut, VOut>(this, redisson, commandExecutor);
}
public int size() {
return get(sizeAsync());
}
public RFuture<Integer> sizeAsync() {
return commandExecutor.readAsync(getRawName(), codec, LLEN_INT, getRawName());
}
public boolean isEmpty() {
return size() == 0;
}
public boolean contains(Object o) {
return get(containsAsync(o));
}
public Iterator<V> iterator() {
return listIterator();
}
public Object[] toArray() {
List<V> list = readAll();
return list.toArray();
}
public List<V> readAll() {
return get(readAllAsync());
}
public RFuture<List<V>> readAllAsync() {
return commandExecutor.readAsync(getRawName(), codec, LRANGE, getRawName(), 0, -1);
}
public <T> T[] toArray(T[] a) {
List<V> list = readAll();
return list.toArray(a);
}
public boolean add(V e) {
return get(addAsync(e));
}
public RFuture<Boolean> addAsync(V e) {
return addAsync(e, RPUSH_BOOLEAN);
}
protected <T> RFuture<T> addAsync(V e, RedisCommand<T> command) {
return commandExecutor.writeAsync(getRawName(), codec, command, getRawName(), encode(e));
}
public boolean remove(Object o) {
return get(removeAsync(o));
}
public RFuture<Boolean> removeAsync(Object o) {
return removeAsync(o, 1);
}
public RFuture<Boolean> removeAsync(Object o, int count) {
return commandExecutor.writeAsync(getRawName(), codec, LREM, getRawName(), count, encode(o));
}
public boolean remove(Object o, int count) {
return get(removeAsync(o, count));
}
public RFuture<Boolean> containsAllAsync(Collection<?> c) {
if (c.isEmpty()) {
return new CompletableFutureWrapper<>(true);
}
return commandExecutor.evalReadAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local items = redis.call('lrange', KEYS[1], 0, -1) " +
"for i=1, #items do " +
"for j = 1, #ARGV, 1 do " +
"if items[i] == ARGV[j] then " +
"table.remove(ARGV, j) " +
"end " +
"end " +
"end " +
"return #ARGV == 0 and 1 or 0",
Collections.<Object>singletonList(getRawName()), encode(c).toArray());
}
public boolean containsAll(Collection<?> c) {
return get(containsAllAsync(c));
}
public boolean addAll(Collection<? extends V> c) {
return get(addAllAsync(c));
}
public RFuture<Boolean> addAllAsync(Collection<? extends V> c) {
if (c.isEmpty()) {
return new CompletableFutureWrapper<>(false);
}
List<Object> args = new ArrayList<Object>(c.size() + 1);
args.add(getRawName());
encode(args, c);
return commandExecutor.writeAsync(getRawName(), codec, RPUSH_BOOLEAN, args.toArray());
}
public RFuture<Boolean> addAllAsync(int index, Collection<? extends V> coll) {
if (index < 0) {
throw new IndexOutOfBoundsException("index: " + index);
}
if (coll.isEmpty()) {
return new CompletableFutureWrapper<>(false);
}
if (index == 0) { // prepend elements to list
List<Object> elements = new ArrayList<Object>();
encode(elements, coll);
Collections.reverse(elements);
elements.add(0, getRawName());
return commandExecutor.writeAsync(getRawName(), codec, LPUSH_BOOLEAN, elements.toArray());
}
List<Object> args = new ArrayList<Object>(coll.size() + 1);
args.add(index);
encode(args, coll);
return commandExecutor.evalWriteNoRetryAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local ind = table.remove(ARGV, 1); " + // index is the first parameter
"local size = redis.call('llen', KEYS[1]); " +
"assert(tonumber(ind) <= size, 'index: ' .. ind .. ' but current size: ' .. size); " +
"local tail = redis.call('lrange', KEYS[1], ind, -1); " +
"redis.call('ltrim', KEYS[1], 0, ind - 1); " +
"for i=1, #ARGV, 5000 do "
+ "redis.call('rpush', KEYS[1], unpack(ARGV, i, math.min(i+4999, #ARGV))); "
+ "end " +
"if #tail > 0 then " +
"for i=1, #tail, 5000 do "
+ "redis.call('rpush', KEYS[1], unpack(tail, i, math.min(i+4999, #tail))); "
+ "end "
+ "end;" +
"return 1;",
Collections.<Object>singletonList(getRawName()), args.toArray());
}
public boolean addAll(int index, Collection<? extends V> coll) {
return get(addAllAsync(index, coll));
}
public RFuture<Boolean> removeAllAsync(Collection<?> c) {
if (c.isEmpty()) {
return new CompletableFutureWrapper<>(false);
}
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local v = 0 " +
"for i = 1, #ARGV, 1 do "
+ "if redis.call('lrem', KEYS[1], 0, ARGV[i]) == 1 "
+ "then v = 1 end "
+"end "
+ "return v ",
Collections.<Object>singletonList(getRawName()), encode(c).toArray());
}
public boolean removeAll(Collection<?> c) {
return get(removeAllAsync(c));
}
public boolean retainAll(Collection<?> c) {
return get(retainAllAsync(c));
}
public RFuture<Boolean> retainAllAsync(Collection<?> c) {
if (c.isEmpty()) {
return deleteAsync();
}
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local changed = 0 " +
"local items = redis.call('lrange', KEYS[1], 0, -1) "
+ "local i = 1 "
+ "while i <= #items do "
+ "local element = items[i] "
+ "local isInAgrs = false "
+ "for j = 1, #ARGV, 1 do "
+ "if ARGV[j] == element then "
+ "isInAgrs = true "
+ "break "
+ "end "
+ "end "
+ "if isInAgrs == false then "
+ "redis.call('LREM', KEYS[1], 0, element) "
+ "changed = 1 "
+ "end "
+ "i = i + 1 "
+ "end "
+ "return changed ",
Collections.<Object>singletonList(getRawName()), encode(c).toArray());
}
public void clear() {
delete();
}
public RFuture<V> getAsync(int index) {
return commandExecutor.readAsync(getRawName(), codec, LINDEX, getRawName(), index);
}
public List<V> get(int... indexes) {
return get(getAsync(indexes));
}
public Iterator<V> distributedIterator(final int count) {
String iteratorName = "__redisson_list_cursor_{" + getRawName() + "}";
return distributedIterator(iteratorName, count);
}
public Iterator<V> distributedIterator(final String iteratorName, final int count) {
return new RedissonBaseIterator<V>() {
@Override
protected ScanResult<Object> iterator(RedisClient client, long nextIterPos) {
return distributedScanIterator(iteratorName, count);
}
@Override
protected void remove(Object value) {
BaseRedissonList.this.remove((V) value);
}
};
}
private ScanResult<Object> distributedScanIterator(String iteratorName, int count) {
return get(distributedScanIteratorAsync(iteratorName, count));
}
private RFuture<ScanResult<Object>> distributedScanIteratorAsync(String iteratorName, int count) {
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_SCAN,
"local start_index = redis.call('get', KEYS[2]); "
+ "if start_index ~= false then "
+ "start_index = tonumber(start_index); "
+ "else "
+ "start_index = 0;"
+ "end;"
+ "if start_index == -1 then "
+ "return {0, {}};"
+ "end;"
+ "local end_index = start_index + ARGV[1];"
+ "local result; "
+ "result = redis.call('lrange', KEYS[1], start_index, end_index - 1); "
+ "if end_index > redis.call('llen', KEYS[1]) then "
+ "end_index = -1;"
+ "end; "
+ "redis.call('setex', KEYS[2], 3600, end_index);"
+ "return {end_index, result};",
Arrays.<Object>asList(getRawName(), iteratorName), count);
}
public RFuture<List<V>> getAsync(int... indexes) {
List<Integer> params = new ArrayList<Integer>();
for (Integer index : indexes) {
params.add(index);
}
return commandExecutor.evalReadAsync(getRawName(), codec, RedisCommands.EVAL_LIST,
"local result = {}; " +
"for i = 1, #ARGV, 1 do "
+ "local value = redis.call('lindex', KEYS[1], ARGV[i]);"
+ "table.insert(result, value);" +
"end; " +
"return result;",
Collections.<Object>singletonList(getRawName()), params.toArray());
}
public V get(int index) {
return getValue(index);
}
V getValue(int index) {
return get(getAsync(index));
}
public V set(int index, V element) {
try {
return get(setAsync(index, element));
} catch (RedisException e) {
if (e.getCause() instanceof IndexOutOfBoundsException) {
throw (IndexOutOfBoundsException) e.getCause();
}
throw e;
}
}
public RFuture<V> setAsync(int index, V element) {
RFuture<V> future = commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_OBJECT,
"local v = redis.call('lindex', KEYS[1], ARGV[1]); " +
"redis.call('lset', KEYS[1], ARGV[1], ARGV[2]); " +
"return v",
Collections.singletonList(getRawName()), index, encode(element));
CompletionStage<V> f = future.handle((res, e) -> {
if (e != null) {
if (e.getMessage().contains("ERR index out of range")) {
throw new CompletionException(new IndexOutOfBoundsException("index out of range"));
}
throw new CompletionException(e);
}
return res;
});
return new CompletableFutureWrapper<>(f);
}
public void fastSet(int index, V element) {
get(fastSetAsync(index, element));
}
public RFuture<Void> fastSetAsync(int index, V element) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LSET, getRawName(), index, encode(element));
}
public void add(int index, V element) {
addAll(index, Collections.singleton(element));
}
public RFuture<Boolean> addAsync(int index, V element) {
return addAllAsync(index, Collections.singleton(element));
}
public V remove(int index) {
return get(removeAsync(index));
}
public RFuture<V> removeAsync(int index) {
if (index == 0) {
return commandExecutor.writeAsync(getRawName(), codec, LPOP, getRawName());
}
return commandExecutor.evalWriteAsync(getRawName(), codec, EVAL_OBJECT,
"local v = redis.call('lindex', KEYS[1], ARGV[1]); " +
"redis.call('lset', KEYS[1], ARGV[1], 'DELETED_BY_REDISSON');" +
"redis.call('lrem', KEYS[1], 1, 'DELETED_BY_REDISSON');" +
"return v",
Collections.<Object>singletonList(getRawName()), index);
}
public void fastRemove(int index) {
get(fastRemoveAsync(index));
}
public RFuture<Void> fastRemoveAsync(int index) {
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_VOID,
"redis.call('lset', KEYS[1], ARGV[1], 'DELETED_BY_REDISSON');" +
"redis.call('lrem', KEYS[1], 1, 'DELETED_BY_REDISSON');",
Collections.<Object>singletonList(getRawName()), index);
}
public int indexOf(Object o) {
return get(indexOfAsync(o));
}
public RFuture<Boolean> containsAsync(Object o) {
return indexOfAsync(o, new BooleanNumberReplayConvertor(-1L));
}
public <R> RFuture<R> indexOfAsync(Object o, Convertor<R> convertor) {
return commandExecutor.evalReadAsync(getRawName(), codec, new RedisCommand<R>("EVAL", convertor),
"local key = KEYS[1] " +
"local obj = ARGV[1] " +
"local items = redis.call('lrange', key, 0, -1) " +
"for i=1,#items do " +
"if items[i] == obj then " +
"return i - 1 " +
"end " +
"end " +
"return -1",
Collections.<Object>singletonList(getRawName()), encode(o));
}
public RFuture<Integer> indexOfAsync(Object o) {
return indexOfAsync(o, new IntegerReplayConvertor());
}
public int lastIndexOf(Object o) {
return get(lastIndexOfAsync(o));
}
public RFuture<Integer> lastIndexOfAsync(Object o) {
return commandExecutor.evalReadAsync(getRawName(), codec, RedisCommands.EVAL_INTEGER,
"local key = KEYS[1] " +
"local obj = ARGV[1] " +
"local items = redis.call('lrange', key, 0, -1) " +
"for i = #items, 1, -1 do " +
"if items[i] == obj then " +
"return i - 1 " +
"end " +
"end " +
"return -1",
Collections.<Object>singletonList(getRawName()), encode(o));
}
public <R> RFuture<R> lastIndexOfAsync(Object o, Convertor<R> convertor) {
return commandExecutor.evalReadAsync(getRawName(), codec, new RedisCommand<R>("EVAL", convertor),
"local key = KEYS[1] " +
"local obj = ARGV[1] " +
"local items = redis.call('lrange', key, 0, -1) " +
"for i = #items, 1, -1 do " +
"if items[i] == obj then " +
"return i - 1 " +
"end " +
"end " +
"return -1",
Collections.<Object>singletonList(getRawName()), encode(o));
}
public void trim(int fromIndex, int toIndex) {
get(trimAsync(fromIndex, toIndex));
}
public RFuture<Void> trimAsync(int fromIndex, int toIndex) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LTRIM, getRawName(), fromIndex, toIndex);
}
public ListIterator<V> listIterator() {
return listIterator(0);
}
public ListIterator<V> listIterator(int ind) {
return new RedissonListIterator<V>(ind) {
@Override
public V getValue(int index) {
return BaseRedissonList.this.getValue(index);
}
@Override
public V remove(int index) {
return BaseRedissonList.this.remove(index);
}
@Override
public void fastSet(int index, V value) {
BaseRedissonList.this.fastSet(index, value);
}
@Override
public void add(int index, V value) {
BaseRedissonList.this.add(index, value);
}
};
}
public RList<V> subList(int fromIndex, int toIndex) {
int size = size();
if (fromIndex < 0 || toIndex > size) {
throw new IndexOutOfBoundsException("fromIndex: " + fromIndex + " toIndex: " + toIndex + " size: " + size);
}
if (fromIndex > toIndex) {
throw new IllegalArgumentException("fromIndex: " + fromIndex + " toIndex: " + toIndex);
}
return new RedissonSubList<V>(codec, commandExecutor, getRawName(), fromIndex, toIndex);
}
@Override
@SuppressWarnings("AvoidInlineConditionals")
public String toString() {
Iterator<V> it = iterator();
if (! it.hasNext())
return "[]";
StringBuilder sb = new StringBuilder();
sb.append('[');
for (;;) {
V e = it.next();
sb.append(e == this ? "(this Collection)" : e);
if (! it.hasNext())
return sb.append(']').toString();
sb.append(',').append(' ');
}
}
@Override
@SuppressWarnings("AvoidInlineConditionals")
public boolean equals(Object o) {
if (o == this)
return true;
if (!(o instanceof List))
return false;
Iterator<V> e1 = iterator();
Iterator<?> e2 = ((List<?>) o).iterator();
while (e1.hasNext() && e2.hasNext()) {
V o1 = e1.next();
Object o2 = e2.next();
if (!(o1==null ? o2==null : o1.equals(o2)))
return false;
}
return !(e1.hasNext() || e2.hasNext());
}
@Override
@SuppressWarnings("AvoidInlineConditionals")
public int hashCode() {
int hashCode = 1;
Iterable<V> ii = () -> iterator();
for (V e : ii) {
hashCode = 31*hashCode + (e==null ? 0 : e.hashCode());
}
return hashCode;
}
public RFuture<Integer> addAfterAsync(V elementToFind, V element) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LINSERT_INT, getRawName(), "AFTER", encode(elementToFind), encode(element));
}
public RFuture<Integer> addBeforeAsync(V elementToFind, V element) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LINSERT_INT, getRawName(), "BEFORE", encode(elementToFind), encode(element));
}
public int addAfter(V elementToFind, V element) {
return get(addAfterAsync(elementToFind, element));
}
public int addBefore(V elementToFind, V element) {
return get(addBeforeAsync(elementToFind, element));
}
public List<V> readSort(SortOrder order) {
return get(readSortAsync(order));
}
public RFuture<List<V>> readSortAsync(SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), order);
}
public List<V> readSort(SortOrder order, int offset, int count) {
return get(readSortAsync(order, offset, count));
}
public RFuture<List<V>> readSortAsync(SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "LIMIT", offset, count, order);
}
public List<V> readSort(String byPattern, SortOrder order) {
return get(readSortAsync(byPattern, order));
}
public RFuture<List<V>> readSortAsync(String byPattern, SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, order);
}
public List<V> readSort(String byPattern, SortOrder order, int offset, int count) {
return get(readSortAsync(byPattern, order, offset, count));
}
public RFuture<List<V>> readSortAsync(String byPattern, SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, "LIMIT", offset, count, order);
}
public <T> Collection<T> readSort(String byPattern, List<String> getPatterns, SortOrder order) {
return (Collection<T>) get(readSortAsync(byPattern, getPatterns, order));
}
public <T> RFuture<Collection<T>> readSortAsync(String byPattern, List<String> getPatterns, SortOrder order) {
return readSortAsync(byPattern, getPatterns, order, -1, -1);
}
public <T> Collection<T> readSort(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return (Collection<T>) get(readSortAsync(byPattern, getPatterns, order, offset, count));
}
public <T> RFuture<Collection<T>> readSortAsync(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return readSortAsync(byPattern, getPatterns, order, offset, count, false);
}
public List<V> readSortAlpha(SortOrder order) {
return get(readSortAlphaAsync(order));
}
public RFuture<List<V>> readSortAlphaAsync(SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "ALPHA", order);
}
public List<V> readSortAlpha(SortOrder order, int offset, int count) {
return get(readSortAlphaAsync(order, offset, count));
}
public RFuture<List<V>> readSortAlphaAsync(SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "LIMIT", offset, count, "ALPHA", order);
}
public List<V> readSortAlpha(String byPattern, SortOrder order) {
return get(readSortAlphaAsync(byPattern, order));
}
public RFuture<List<V>> readSortAlphaAsync(String byPattern, SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, "ALPHA", order);
}
public List<V> readSortAlpha(String byPattern, SortOrder order, int offset, int count) {
return get(readSortAlphaAsync(byPattern, order, offset, count));
}
public RFuture<List<V>> readSortAlphaAsync(String byPattern, SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, "LIMIT", offset, count, "ALPHA", order);
}
public <T> Collection<T> readSortAlpha(String byPattern, List<String> getPatterns, SortOrder order) {
return (Collection<T>) get(readSortAlphaAsync(byPattern, getPatterns, order));
}
public <T> RFuture<Collection<T>> readSortAlphaAsync(String byPattern, List<String> getPatterns, SortOrder order) {
return readSortAlphaAsync(byPattern, getPatterns, order, -1, -1);
}
public <T> Collection<T> readSortAlpha(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return (Collection<T>) get(readSortAlphaAsync(byPattern, getPatterns, order, offset, count));
}
public <T> RFuture<Collection<T>> readSortAlphaAsync(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return readSortAsync(byPattern, getPatterns, order, offset, count, true);
}
public int sortTo(String destName, SortOrder order) {
return get(sortToAsync(destName, order));
}
public RFuture<Integer> sortToAsync(String destName, SortOrder order) {
return sortToAsync(destName, null, Collections.<String>emptyList(), order, -1, -1);
}
public int sortTo(String destName, SortOrder order, int offset, int count) {
return get(sortToAsync(destName, order, offset, count));
}
public RFuture<Integer> sortToAsync(String destName, SortOrder order, int offset, int count) {
return sortToAsync(destName, null, Collections.<String>emptyList(), order, offset, count);
}
public int sortTo(String destName, String byPattern, SortOrder order, int offset, int count) {
return get(sortToAsync(destName, byPattern, order, offset, count));
}
public int sortTo(String destName, String byPattern, SortOrder order) {
return get(sortToAsync(destName, byPattern, order));
}
public RFuture<Integer> sortToAsync(String destName, String byPattern, SortOrder order) {
return sortToAsync(destName, byPattern, Collections.<String>emptyList(), order, -1, -1);
}
public RFuture<Integer> sortToAsync(String destName, String byPattern, SortOrder order, int offset, int count) {
return sortToAsync(destName, byPattern, Collections.<String>emptyList(), order, offset, count);
}
public int sortTo(String destName, String byPattern, List<String> getPatterns, SortOrder order) {
return get(sortToAsync(destName, byPattern, getPatterns, order));
}
public RFuture<Integer> sortToAsync(String destName, String byPattern, List<String> getPatterns, SortOrder order) {
return sortToAsync(destName, byPattern, getPatterns, order, -1, -1);
}
public int sortTo(String destName, String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return get(sortToAsync(destName, byPattern, getPatterns, order, offset, count));
}
public RFuture<Integer> sortToAsync(String destName, String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
List<Object> params = new ArrayList<Object>();
params.add(getRawName());
if (byPattern != null) {
params.add("BY");
params.add(byPattern);
}
if (offset != -1 && count != -1) {
params.add("LIMIT");
}
if (offset != -1) {
params.add(offset);
}
if (count != -1) {
params.add(count);
}
for (String pattern : getPatterns) {
params.add("GET");
params.add(pattern);
}
params.add(order);
params.add("STORE");
params.add(destName);
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.SORT_TO, params.toArray());
}
private <T> RFuture<Collection<T>> readSortAsync(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count, boolean alpha) {
List<Object> params = new ArrayList<Object>();
params.add(getRawName());
if (byPattern != null) {
params.add("BY");
params.add(byPattern);
}
if (offset != -1 && count != -1) {
params.add("LIMIT");
}
if (offset != -1) {
params.add(offset);
}
if (count != -1) {
params.add(count);
}
if (getPatterns != null) {
for (String pattern : getPatterns) {
params.add("GET");
params.add(pattern);
}
}
if (alpha) {
params.add("ALPHA");
}
if (order != null) {
params.add(order);
}
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, params.toArray());
}
public RFuture<List<V>> rangeAsync(int toIndex) {
return rangeAsync(0, toIndex);
}
public RFuture<List<V>> rangeAsync(int fromIndex, int toIndex) {
return commandExecutor.readAsync(getRawName(), codec, LRANGE, getRawName(), fromIndex, toIndex);
}
public List<V> range(int toIndex) {
return get(rangeAsync(toIndex));
}
public List<V> range(int fromIndex, int toIndex) {
return get(rangeAsync(fromIndex, toIndex));
}
@Override
public int addListener(ObjectListener listener) {
if (listener instanceof ListAddListener) {
return addListener("__keyevent@*:rpush", (ListAddListener) listener, ListAddListener::onListAdd);
}
if (listener instanceof ListRemoveListener) {
return addListener("__keyevent@*:lrem", (ListRemoveListener) listener, ListRemoveListener::onListRemove);
}
if (listener instanceof ListTrimListener) {
return addListener("__keyevent@*:ltrim", (ListTrimListener) listener, ListTrimListener::onListTrim);
}
if (listener instanceof ListSetListener) {
return addListener("__keyevent@*:lset", (ListSetListener) listener, ListSetListener::onListSet);
}
if (listener instanceof ListInsertListener) {
return addListener("__keyevent@*:linsert", (ListInsertListener) listener, ListInsertListener::onListInsert);
}
return super.addListener(listener);
}
@Override
public RFuture<Integer> addListenerAsync(ObjectListener listener) {
if (listener instanceof ListAddListener) {
return addListenerAsync("__keyevent@*:rpush", (ListAddListener) listener, ListAddListener::onListAdd);
}
if (listener instanceof ListRemoveListener) {
return addListenerAsync("__keyevent@*:lrem", (ListRemoveListener) listener, ListRemoveListener::onListRemove);
}
if (listener instanceof ListTrimListener) {
return addListenerAsync("__keyevent@*:ltrim", (ListTrimListener) listener, ListTrimListener::onListTrim);
}
if (listener instanceof ListSetListener) {
return addListenerAsync("__keyevent@*:lset", (ListSetListener) listener, ListSetListener::onListSet);
}
if (listener instanceof ListInsertListener) {
return addListenerAsync("__keyevent@*:linsert", (ListInsertListener) listener, ListInsertListener::onListInsert);
}
return super.addListenerAsync(listener);
}
@Override
public void removeListener(int listenerId) {
RPatternTopic addTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:rpush");
addTopic.removeListener(listenerId);
RPatternTopic remTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lrem");
remTopic.removeListener(listenerId);
RPatternTopic trimTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:ltrim");
trimTopic.removeListener(listenerId);
RPatternTopic setTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lset");
setTopic.removeListener(listenerId);
RPatternTopic insertTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:linsert");
insertTopic.removeListener(listenerId);
super.removeListener(listenerId);
}
@Override
public RFuture<Void> removeListenerAsync(int listenerId) {
RPatternTopic addTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:rpush");
RFuture<Void> f1 = addTopic.removeListenerAsync(listenerId);
RPatternTopic remTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lrem");
RFuture<Void> f2 = remTopic.removeListenerAsync(listenerId);
RPatternTopic trimTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:ltrim");
RFuture<Void> f3 = trimTopic.removeListenerAsync(listenerId);
RPatternTopic setTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lset");
RFuture<Void> f4 = setTopic.removeListenerAsync(listenerId);
RPatternTopic insertTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:linsert");
RFuture<Void> f5 = insertTopic.removeListenerAsync(listenerId);
RFuture<Void> f6 = super.removeListenerAsync(listenerId);
CompletableFuture<Void> f = CompletableFuture.allOf(f1.toCompletableFuture(), f2.toCompletableFuture(), f3.toCompletableFuture(),
f4.toCompletableFuture(), f5.toCompletableFuture(), f5.toCompletableFuture(), f6.toCompletableFuture());
return new CompletableFutureWrapper<>(f);
}
public boolean removeIf(Predicate<? super V> filter) {
throw new UnsupportedOperationException();
}
}

@ -15,8 +15,6 @@
*/
package org.redisson;
import java.util.*;
import org.redisson.api.RDeque;
import org.redisson.api.RFuture;
import org.redisson.api.RedissonClient;
@ -28,6 +26,8 @@ import org.redisson.client.protocol.RedisCommands;
import org.redisson.client.protocol.decoder.ListFirstObjectDecoder;
import org.redisson.command.CommandAsyncExecutor;
import java.util.*;
/**
* Distributed and concurrent implementation of {@link java.util.Queue}
*
@ -332,8 +332,4 @@ public class RedissonDeque<V> extends RedissonQueue<V> implements RDeque<V> {
return remove(o, -1);
}
public RedissonDeque<V> reversed() {
throw new UnsupportedOperationException();
}
}

@ -22,11 +22,18 @@ import org.redisson.client.codec.LongCodec;
import org.redisson.client.codec.StringCodec;
import org.redisson.client.protocol.RedisCommand;
import org.redisson.client.protocol.RedisCommands;
import org.redisson.client.protocol.RedisStrictCommand;
import org.redisson.client.protocol.convertor.JsonTypeConvertor;
import org.redisson.client.protocol.convertor.LongNumberConvertor;
import org.redisson.client.protocol.convertor.NumberConvertor;
import org.redisson.client.protocol.decoder.ListFirstObjectDecoder;
import org.redisson.client.protocol.decoder.ListMultiDecoder2;
import org.redisson.client.protocol.decoder.ObjectListReplayDecoder;
import org.redisson.client.protocol.decoder.StringListListReplayDecoder;
import org.redisson.codec.JsonCodec;
import org.redisson.codec.JsonCodecWrapper;
import org.redisson.command.CommandAsyncExecutor;
import org.redisson.config.Protocol;
import java.math.BigDecimal;
import java.time.Duration;
@ -87,6 +94,9 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public RFuture<V> getAsync() {
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.JSON_GET, getRawName(), ".");
}
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.JSON_GET, getRawName());
}
@ -97,6 +107,12 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public <T> RFuture<T> getAsync(JsonCodec<T> codec, String... paths) {
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
if (paths.length == 0) {
paths = new String[]{"."};
}
}
List<Object> args = new ArrayList<>();
args.add(getRawName());
args.addAll(Arrays.asList(paths));
@ -761,7 +777,13 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public <T extends Number> RFuture<T> incrementAndGetAsync(String path, T delta) {
return commandExecutor.writeAsync(getRawName(), StringCodec.INSTANCE, new RedisCommand<>("JSON.NUMINCRBY", new NumberConvertor(delta.getClass())),
RedisCommand command;
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
command = new RedisCommand<>("JSON.NUMINCRBY", new ListFirstObjectDecoder(), new LongNumberConvertor(delta.getClass()));
} else {
command = new RedisCommand<>("JSON.NUMINCRBY", new NumberConvertor(delta.getClass()));
}
return commandExecutor.writeAsync(getRawName(), StringCodec.INSTANCE, command,
getRawName(), path, new BigDecimal(delta.toString()).toPlainString());
}
@ -784,7 +806,12 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public RFuture<Long> countKeysAsync() {
return commandExecutor.writeAsync(getRawName(), LongCodec.INSTANCE, RedisCommands.JSON_OBJLEN, getRawName());
RedisStrictCommand command = RedisCommands.JSON_OBJLEN;
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
command = new RedisStrictCommand("JSON.OBJLEN", new ListFirstObjectDecoder());
}
return commandExecutor.writeAsync(getRawName(), LongCodec.INSTANCE, command, getRawName());
}
@Override
@ -814,7 +841,12 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public RFuture<List<String>> getKeysAsync() {
return commandExecutor.readAsync(getRawName(), LongCodec.INSTANCE, RedisCommands.JSON_OBJKEYS, getRawName());
RedisCommand command = RedisCommands.JSON_OBJKEYS;
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
command = new RedisCommand("JSON.OBJKEYS",
new ListMultiDecoder2(new ListFirstObjectDecoder(), new StringListListReplayDecoder()));
}
return commandExecutor.readAsync(getRawName(), LongCodec.INSTANCE, command, getRawName());
}
@Override
@ -864,7 +896,12 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public RFuture<JsonType> getTypeAsync() {
return commandExecutor.readAsync(getRawName(), StringCodec.INSTANCE, RedisCommands.JSON_TYPE, getRawName());
RedisCommand command = RedisCommands.JSON_TYPE;
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
command = new RedisCommand("JSON.TYPE", new ListFirstObjectDecoder(), new JsonTypeConvertor());
}
return commandExecutor.readAsync(getRawName(), StringCodec.INSTANCE, command, getRawName());
}
@Override
@ -874,7 +911,12 @@ public class RedissonJsonBucket<V> extends RedissonExpirable implements RJsonBuc
@Override
public RFuture<JsonType> getTypeAsync(String path) {
return commandExecutor.readAsync(getRawName(), StringCodec.INSTANCE, RedisCommands.JSON_TYPE, getRawName(), path);
RedisCommand command = RedisCommands.JSON_TYPE;
if (getServiceManager().getCfg().getProtocol() == Protocol.RESP3) {
command = new RedisCommand("JSON.TYPE", new ListFirstObjectDecoder(), new JsonTypeConvertor());
}
return commandExecutor.readAsync(getRawName(), StringCodec.INSTANCE, command, getRawName(), path);
}
@Override

@ -15,31 +15,10 @@
*/
package org.redisson;
import org.redisson.api.*;
import org.redisson.api.listener.*;
import org.redisson.api.mapreduce.RCollectionMapReduce;
import org.redisson.client.RedisClient;
import org.redisson.client.RedisException;
import org.redisson.api.RList;
import org.redisson.api.RedissonClient;
import org.redisson.client.codec.Codec;
import org.redisson.client.codec.StringCodec;
import org.redisson.client.protocol.RedisCommand;
import org.redisson.client.protocol.RedisCommands;
import org.redisson.client.protocol.convertor.BooleanNumberReplayConvertor;
import org.redisson.client.protocol.convertor.Convertor;
import org.redisson.client.protocol.convertor.IntegerReplayConvertor;
import org.redisson.command.CommandAsyncExecutor;
import org.redisson.iterator.RedissonBaseIterator;
import org.redisson.iterator.RedissonListIterator;
import org.redisson.mapreduce.RedissonCollectionMapReduce;
import org.redisson.misc.CompletableFutureWrapper;
import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.CompletionStage;
import java.util.function.Predicate;
import static org.redisson.client.protocol.RedisCommands.*;
/**
* Distributed and concurrent implementation of {@link java.util.List}
@ -48,944 +27,14 @@ import static org.redisson.client.protocol.RedisCommands.*;
*
* @param <V> the type of elements held in this collection
*/
public class RedissonList<V> extends RedissonExpirable implements RList<V> {
public class RedissonList<V> extends BaseRedissonList<V> implements RList<V> {
private RedissonClient redisson;
public RedissonList(CommandAsyncExecutor commandExecutor, String name, RedissonClient redisson) {
super(commandExecutor, name);
this.redisson = redisson;
super(commandExecutor, name, redisson);
}
public RedissonList(Codec codec, CommandAsyncExecutor commandExecutor, String name, RedissonClient redisson) {
super(codec, commandExecutor, name);
this.redisson = redisson;
}
@Override
public <KOut, VOut> RCollectionMapReduce<V, KOut, VOut> mapReduce() {
return new RedissonCollectionMapReduce<V, KOut, VOut>(this, redisson, commandExecutor);
}
@Override
public int size() {
return get(sizeAsync());
}
public RFuture<Integer> sizeAsync() {
return commandExecutor.readAsync(getRawName(), codec, LLEN_INT, getRawName());
}
@Override
public boolean isEmpty() {
return size() == 0;
}
@Override
public boolean contains(Object o) {
return get(containsAsync(o));
}
@Override
public Iterator<V> iterator() {
return listIterator();
}
@Override
public Object[] toArray() {
List<V> list = readAll();
return list.toArray();
}
@Override
public List<V> readAll() {
return get(readAllAsync());
}
@Override
public RFuture<List<V>> readAllAsync() {
return commandExecutor.readAsync(getRawName(), codec, LRANGE, getRawName(), 0, -1);
}
@Override
public <T> T[] toArray(T[] a) {
List<V> list = readAll();
return list.toArray(a);
}
@Override
public boolean add(V e) {
return get(addAsync(e));
}
@Override
public RFuture<Boolean> addAsync(V e) {
return addAsync(e, RPUSH_BOOLEAN);
}
protected <T> RFuture<T> addAsync(V e, RedisCommand<T> command) {
return commandExecutor.writeAsync(getRawName(), codec, command, getRawName(), encode(e));
}
@Override
public boolean remove(Object o) {
return get(removeAsync(o));
}
@Override
public RFuture<Boolean> removeAsync(Object o) {
return removeAsync(o, 1);
}
@Override
public RFuture<Boolean> removeAsync(Object o, int count) {
return commandExecutor.writeAsync(getRawName(), codec, LREM, getRawName(), count, encode(o));
}
@Override
public boolean remove(Object o, int count) {
return get(removeAsync(o, count));
}
@Override
public RFuture<Boolean> containsAllAsync(Collection<?> c) {
if (c.isEmpty()) {
return new CompletableFutureWrapper<>(true);
}
return commandExecutor.evalReadAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local items = redis.call('lrange', KEYS[1], 0, -1) " +
"for i=1, #items do " +
"for j = 1, #ARGV, 1 do " +
"if items[i] == ARGV[j] then " +
"table.remove(ARGV, j) " +
"end " +
"end " +
"end " +
"return #ARGV == 0 and 1 or 0",
Collections.<Object>singletonList(getRawName()), encode(c).toArray());
}
@Override
public boolean containsAll(Collection<?> c) {
return get(containsAllAsync(c));
}
@Override
public boolean addAll(Collection<? extends V> c) {
return get(addAllAsync(c));
}
@Override
public RFuture<Boolean> addAllAsync(Collection<? extends V> c) {
if (c.isEmpty()) {
return new CompletableFutureWrapper<>(false);
}
List<Object> args = new ArrayList<Object>(c.size() + 1);
args.add(getRawName());
encode(args, c);
return commandExecutor.writeAsync(getRawName(), codec, RPUSH_BOOLEAN, args.toArray());
}
@Override
public RFuture<Boolean> addAllAsync(int index, Collection<? extends V> coll) {
if (index < 0) {
throw new IndexOutOfBoundsException("index: " + index);
}
if (coll.isEmpty()) {
return new CompletableFutureWrapper<>(false);
}
if (index == 0) { // prepend elements to list
List<Object> elements = new ArrayList<Object>();
encode(elements, coll);
Collections.reverse(elements);
elements.add(0, getRawName());
return commandExecutor.writeAsync(getRawName(), codec, LPUSH_BOOLEAN, elements.toArray());
}
List<Object> args = new ArrayList<Object>(coll.size() + 1);
args.add(index);
encode(args, coll);
return commandExecutor.evalWriteNoRetryAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local ind = table.remove(ARGV, 1); " + // index is the first parameter
"local size = redis.call('llen', KEYS[1]); " +
"assert(tonumber(ind) <= size, 'index: ' .. ind .. ' but current size: ' .. size); " +
"local tail = redis.call('lrange', KEYS[1], ind, -1); " +
"redis.call('ltrim', KEYS[1], 0, ind - 1); " +
"for i=1, #ARGV, 5000 do "
+ "redis.call('rpush', KEYS[1], unpack(ARGV, i, math.min(i+4999, #ARGV))); "
+ "end " +
"if #tail > 0 then " +
"for i=1, #tail, 5000 do "
+ "redis.call('rpush', KEYS[1], unpack(tail, i, math.min(i+4999, #tail))); "
+ "end "
+ "end;" +
"return 1;",
Collections.<Object>singletonList(getRawName()), args.toArray());
}
@Override
public boolean addAll(int index, Collection<? extends V> coll) {
return get(addAllAsync(index, coll));
}
@Override
public RFuture<Boolean> removeAllAsync(Collection<?> c) {
if (c.isEmpty()) {
return new CompletableFutureWrapper<>(false);
}
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local v = 0 " +
"for i = 1, #ARGV, 1 do "
+ "if redis.call('lrem', KEYS[1], 0, ARGV[i]) == 1 "
+ "then v = 1 end "
+"end "
+ "return v ",
Collections.<Object>singletonList(getRawName()), encode(c).toArray());
}
@Override
public boolean removeAll(Collection<?> c) {
return get(removeAllAsync(c));
}
@Override
public boolean retainAll(Collection<?> c) {
return get(retainAllAsync(c));
}
@Override
public RFuture<Boolean> retainAllAsync(Collection<?> c) {
if (c.isEmpty()) {
return deleteAsync();
}
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_BOOLEAN,
"local changed = 0 " +
"local items = redis.call('lrange', KEYS[1], 0, -1) "
+ "local i = 1 "
+ "while i <= #items do "
+ "local element = items[i] "
+ "local isInAgrs = false "
+ "for j = 1, #ARGV, 1 do "
+ "if ARGV[j] == element then "
+ "isInAgrs = true "
+ "break "
+ "end "
+ "end "
+ "if isInAgrs == false then "
+ "redis.call('LREM', KEYS[1], 0, element) "
+ "changed = 1 "
+ "end "
+ "i = i + 1 "
+ "end "
+ "return changed ",
Collections.<Object>singletonList(getRawName()), encode(c).toArray());
}
@Override
public void clear() {
delete();
}
@Override
public RFuture<V> getAsync(int index) {
return commandExecutor.readAsync(getRawName(), codec, LINDEX, getRawName(), index);
}
public List<V> get(int... indexes) {
return get(getAsync(indexes));
}
@Override
public Iterator<V> distributedIterator(final int count) {
String iteratorName = "__redisson_list_cursor_{" + getRawName() + "}";
return distributedIterator(iteratorName, count);
}
@Override
public Iterator<V> distributedIterator(final String iteratorName, final int count) {
return new RedissonBaseIterator<V>() {
@Override
protected ScanResult<Object> iterator(RedisClient client, long nextIterPos) {
return distributedScanIterator(iteratorName, count);
}
@Override
protected void remove(Object value) {
RedissonList.this.remove((V) value);
}
};
}
private ScanResult<Object> distributedScanIterator(String iteratorName, int count) {
return get(distributedScanIteratorAsync(iteratorName, count));
}
private RFuture<ScanResult<Object>> distributedScanIteratorAsync(String iteratorName, int count) {
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_SCAN,
"local start_index = redis.call('get', KEYS[2]); "
+ "if start_index ~= false then "
+ "start_index = tonumber(start_index); "
+ "else "
+ "start_index = 0;"
+ "end;"
+ "if start_index == -1 then "
+ "return {0, {}};"
+ "end;"
+ "local end_index = start_index + ARGV[1];"
+ "local result; "
+ "result = redis.call('lrange', KEYS[1], start_index, end_index - 1); "
+ "if end_index > redis.call('llen', KEYS[1]) then "
+ "end_index = -1;"
+ "end; "
+ "redis.call('setex', KEYS[2], 3600, end_index);"
+ "return {end_index, result};",
Arrays.<Object>asList(getRawName(), iteratorName), count);
}
public RFuture<List<V>> getAsync(int... indexes) {
List<Integer> params = new ArrayList<Integer>();
for (Integer index : indexes) {
params.add(index);
}
return commandExecutor.evalReadAsync(getRawName(), codec, RedisCommands.EVAL_LIST,
"local result = {}; " +
"for i = 1, #ARGV, 1 do "
+ "local value = redis.call('lindex', KEYS[1], ARGV[i]);"
+ "table.insert(result, value);" +
"end; " +
"return result;",
Collections.<Object>singletonList(getRawName()), params.toArray());
}
@Override
public V get(int index) {
return getValue(index);
}
V getValue(int index) {
return get(getAsync(index));
}
@Override
public V set(int index, V element) {
try {
return get(setAsync(index, element));
} catch (RedisException e) {
if (e.getCause() instanceof IndexOutOfBoundsException) {
throw (IndexOutOfBoundsException) e.getCause();
}
throw e;
}
}
@Override
public RFuture<V> setAsync(int index, V element) {
RFuture<V> future = commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_OBJECT,
"local v = redis.call('lindex', KEYS[1], ARGV[1]); " +
"redis.call('lset', KEYS[1], ARGV[1], ARGV[2]); " +
"return v",
Collections.singletonList(getRawName()), index, encode(element));
CompletionStage<V> f = future.handle((res, e) -> {
if (e != null) {
if (e.getMessage().contains("ERR index out of range")) {
throw new CompletionException(new IndexOutOfBoundsException("index out of range"));
}
throw new CompletionException(e);
}
return res;
});
return new CompletableFutureWrapper<>(f);
}
@Override
public void fastSet(int index, V element) {
get(fastSetAsync(index, element));
}
@Override
public RFuture<Void> fastSetAsync(int index, V element) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LSET, getRawName(), index, encode(element));
}
@Override
public void add(int index, V element) {
addAll(index, Collections.singleton(element));
}
@Override
public RFuture<Boolean> addAsync(int index, V element) {
return addAllAsync(index, Collections.singleton(element));
}
@Override
public V remove(int index) {
return get(removeAsync(index));
}
@Override
public RFuture<V> removeAsync(int index) {
if (index == 0) {
return commandExecutor.writeAsync(getRawName(), codec, LPOP, getRawName());
}
return commandExecutor.evalWriteAsync(getRawName(), codec, EVAL_OBJECT,
"local v = redis.call('lindex', KEYS[1], ARGV[1]); " +
"redis.call('lset', KEYS[1], ARGV[1], 'DELETED_BY_REDISSON');" +
"redis.call('lrem', KEYS[1], 1, 'DELETED_BY_REDISSON');" +
"return v",
Collections.<Object>singletonList(getRawName()), index);
}
@Override
public void fastRemove(int index) {
get(fastRemoveAsync(index));
}
@Override
public RFuture<Void> fastRemoveAsync(int index) {
return commandExecutor.evalWriteAsync(getRawName(), codec, RedisCommands.EVAL_VOID,
"redis.call('lset', KEYS[1], ARGV[1], 'DELETED_BY_REDISSON');" +
"redis.call('lrem', KEYS[1], 1, 'DELETED_BY_REDISSON');",
Collections.<Object>singletonList(getRawName()), index);
}
@Override
public int indexOf(Object o) {
return get(indexOfAsync(o));
}
@Override
public RFuture<Boolean> containsAsync(Object o) {
return indexOfAsync(o, new BooleanNumberReplayConvertor(-1L));
}
public <R> RFuture<R> indexOfAsync(Object o, Convertor<R> convertor) {
return commandExecutor.evalReadAsync(getRawName(), codec, new RedisCommand<R>("EVAL", convertor),
"local key = KEYS[1] " +
"local obj = ARGV[1] " +
"local items = redis.call('lrange', key, 0, -1) " +
"for i=1,#items do " +
"if items[i] == obj then " +
"return i - 1 " +
"end " +
"end " +
"return -1",
Collections.<Object>singletonList(getRawName()), encode(o));
}
@Override
public RFuture<Integer> indexOfAsync(Object o) {
return indexOfAsync(o, new IntegerReplayConvertor());
}
@Override
public int lastIndexOf(Object o) {
return get(lastIndexOfAsync(o));
}
@Override
public RFuture<Integer> lastIndexOfAsync(Object o) {
return commandExecutor.evalReadAsync(getRawName(), codec, RedisCommands.EVAL_INTEGER,
"local key = KEYS[1] " +
"local obj = ARGV[1] " +
"local items = redis.call('lrange', key, 0, -1) " +
"for i = #items, 1, -1 do " +
"if items[i] == obj then " +
"return i - 1 " +
"end " +
"end " +
"return -1",
Collections.<Object>singletonList(getRawName()), encode(o));
}
public <R> RFuture<R> lastIndexOfAsync(Object o, Convertor<R> convertor) {
return commandExecutor.evalReadAsync(getRawName(), codec, new RedisCommand<R>("EVAL", convertor),
"local key = KEYS[1] " +
"local obj = ARGV[1] " +
"local items = redis.call('lrange', key, 0, -1) " +
"for i = #items, 1, -1 do " +
"if items[i] == obj then " +
"return i - 1 " +
"end " +
"end " +
"return -1",
Collections.<Object>singletonList(getRawName()), encode(o));
}
@Override
public void trim(int fromIndex, int toIndex) {
get(trimAsync(fromIndex, toIndex));
}
@Override
public RFuture<Void> trimAsync(int fromIndex, int toIndex) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LTRIM, getRawName(), fromIndex, toIndex);
super(codec, commandExecutor, name, redisson);
}
@Override
public ListIterator<V> listIterator() {
return listIterator(0);
}
@Override
public ListIterator<V> listIterator(int ind) {
return new RedissonListIterator<V>(ind) {
@Override
public V getValue(int index) {
return RedissonList.this.getValue(index);
}
@Override
public V remove(int index) {
return RedissonList.this.remove(index);
}
@Override
public void fastSet(int index, V value) {
RedissonList.this.fastSet(index, value);
}
@Override
public void add(int index, V value) {
RedissonList.this.add(index, value);
}
};
}
@Override
public RList<V> subList(int fromIndex, int toIndex) {
int size = size();
if (fromIndex < 0 || toIndex > size) {
throw new IndexOutOfBoundsException("fromIndex: " + fromIndex + " toIndex: " + toIndex + " size: " + size);
}
if (fromIndex > toIndex) {
throw new IllegalArgumentException("fromIndex: " + fromIndex + " toIndex: " + toIndex);
}
return new RedissonSubList<V>(codec, commandExecutor, getRawName(), fromIndex, toIndex);
}
@Override
@SuppressWarnings("AvoidInlineConditionals")
public String toString() {
Iterator<V> it = iterator();
if (! it.hasNext())
return "[]";
StringBuilder sb = new StringBuilder();
sb.append('[');
for (;;) {
V e = it.next();
sb.append(e == this ? "(this Collection)" : e);
if (! it.hasNext())
return sb.append(']').toString();
sb.append(',').append(' ');
}
}
@Override
@SuppressWarnings("AvoidInlineConditionals")
public boolean equals(Object o) {
if (o == this)
return true;
if (!(o instanceof List))
return false;
Iterator<V> e1 = iterator();
Iterator<?> e2 = ((List<?>) o).iterator();
while (e1.hasNext() && e2.hasNext()) {
V o1 = e1.next();
Object o2 = e2.next();
if (!(o1==null ? o2==null : o1.equals(o2)))
return false;
}
return !(e1.hasNext() || e2.hasNext());
}
@Override
@SuppressWarnings("AvoidInlineConditionals")
public int hashCode() {
int hashCode = 1;
for (V e : this) {
hashCode = 31*hashCode + (e==null ? 0 : e.hashCode());
}
return hashCode;
}
@Override
public RFuture<Integer> addAfterAsync(V elementToFind, V element) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LINSERT_INT, getRawName(), "AFTER", encode(elementToFind), encode(element));
}
@Override
public RFuture<Integer> addBeforeAsync(V elementToFind, V element) {
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.LINSERT_INT, getRawName(), "BEFORE", encode(elementToFind), encode(element));
}
@Override
public int addAfter(V elementToFind, V element) {
return get(addAfterAsync(elementToFind, element));
}
@Override
public int addBefore(V elementToFind, V element) {
return get(addBeforeAsync(elementToFind, element));
}
@Override
public List<V> readSort(SortOrder order) {
return get(readSortAsync(order));
}
@Override
public RFuture<List<V>> readSortAsync(SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), order);
}
@Override
public List<V> readSort(SortOrder order, int offset, int count) {
return get(readSortAsync(order, offset, count));
}
@Override
public RFuture<List<V>> readSortAsync(SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "LIMIT", offset, count, order);
}
@Override
public List<V> readSort(String byPattern, SortOrder order) {
return get(readSortAsync(byPattern, order));
}
@Override
public RFuture<List<V>> readSortAsync(String byPattern, SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, order);
}
@Override
public List<V> readSort(String byPattern, SortOrder order, int offset, int count) {
return get(readSortAsync(byPattern, order, offset, count));
}
@Override
public RFuture<List<V>> readSortAsync(String byPattern, SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, "LIMIT", offset, count, order);
}
@Override
public <T> Collection<T> readSort(String byPattern, List<String> getPatterns, SortOrder order) {
return (Collection<T>) get(readSortAsync(byPattern, getPatterns, order));
}
@Override
public <T> RFuture<Collection<T>> readSortAsync(String byPattern, List<String> getPatterns, SortOrder order) {
return readSortAsync(byPattern, getPatterns, order, -1, -1);
}
@Override
public <T> Collection<T> readSort(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return (Collection<T>) get(readSortAsync(byPattern, getPatterns, order, offset, count));
}
@Override
public <T> RFuture<Collection<T>> readSortAsync(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return readSortAsync(byPattern, getPatterns, order, offset, count, false);
}
@Override
public List<V> readSortAlpha(SortOrder order) {
return get(readSortAlphaAsync(order));
}
@Override
public RFuture<List<V>> readSortAlphaAsync(SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "ALPHA", order);
}
@Override
public List<V> readSortAlpha(SortOrder order, int offset, int count) {
return get(readSortAlphaAsync(order, offset, count));
}
@Override
public RFuture<List<V>> readSortAlphaAsync(SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "LIMIT", offset, count, "ALPHA", order);
}
@Override
public List<V> readSortAlpha(String byPattern, SortOrder order) {
return get(readSortAlphaAsync(byPattern, order));
}
@Override
public RFuture<List<V>> readSortAlphaAsync(String byPattern, SortOrder order) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, "ALPHA", order);
}
@Override
public List<V> readSortAlpha(String byPattern, SortOrder order, int offset, int count) {
return get(readSortAlphaAsync(byPattern, order, offset, count));
}
@Override
public RFuture<List<V>> readSortAlphaAsync(String byPattern, SortOrder order, int offset, int count) {
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, getRawName(), "BY", byPattern, "LIMIT", offset, count, "ALPHA", order);
}
@Override
public <T> Collection<T> readSortAlpha(String byPattern, List<String> getPatterns, SortOrder order) {
return (Collection<T>) get(readSortAlphaAsync(byPattern, getPatterns, order));
}
@Override
public <T> RFuture<Collection<T>> readSortAlphaAsync(String byPattern, List<String> getPatterns, SortOrder order) {
return readSortAlphaAsync(byPattern, getPatterns, order, -1, -1);
}
@Override
public <T> Collection<T> readSortAlpha(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return (Collection<T>) get(readSortAlphaAsync(byPattern, getPatterns, order, offset, count));
}
@Override
public <T> RFuture<Collection<T>> readSortAlphaAsync(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return readSortAsync(byPattern, getPatterns, order, offset, count, true);
}
@Override
public int sortTo(String destName, SortOrder order) {
return get(sortToAsync(destName, order));
}
@Override
public RFuture<Integer> sortToAsync(String destName, SortOrder order) {
return sortToAsync(destName, null, Collections.<String>emptyList(), order, -1, -1);
}
@Override
public int sortTo(String destName, SortOrder order, int offset, int count) {
return get(sortToAsync(destName, order, offset, count));
}
@Override
public RFuture<Integer> sortToAsync(String destName, SortOrder order, int offset, int count) {
return sortToAsync(destName, null, Collections.<String>emptyList(), order, offset, count);
}
@Override
public int sortTo(String destName, String byPattern, SortOrder order, int offset, int count) {
return get(sortToAsync(destName, byPattern, order, offset, count));
}
@Override
public int sortTo(String destName, String byPattern, SortOrder order) {
return get(sortToAsync(destName, byPattern, order));
}
@Override
public RFuture<Integer> sortToAsync(String destName, String byPattern, SortOrder order) {
return sortToAsync(destName, byPattern, Collections.<String>emptyList(), order, -1, -1);
}
@Override
public RFuture<Integer> sortToAsync(String destName, String byPattern, SortOrder order, int offset, int count) {
return sortToAsync(destName, byPattern, Collections.<String>emptyList(), order, offset, count);
}
@Override
public int sortTo(String destName, String byPattern, List<String> getPatterns, SortOrder order) {
return get(sortToAsync(destName, byPattern, getPatterns, order));
}
@Override
public RFuture<Integer> sortToAsync(String destName, String byPattern, List<String> getPatterns, SortOrder order) {
return sortToAsync(destName, byPattern, getPatterns, order, -1, -1);
}
@Override
public int sortTo(String destName, String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
return get(sortToAsync(destName, byPattern, getPatterns, order, offset, count));
}
@Override
public RFuture<Integer> sortToAsync(String destName, String byPattern, List<String> getPatterns, SortOrder order, int offset, int count) {
List<Object> params = new ArrayList<Object>();
params.add(getRawName());
if (byPattern != null) {
params.add("BY");
params.add(byPattern);
}
if (offset != -1 && count != -1) {
params.add("LIMIT");
}
if (offset != -1) {
params.add(offset);
}
if (count != -1) {
params.add(count);
}
for (String pattern : getPatterns) {
params.add("GET");
params.add(pattern);
}
params.add(order);
params.add("STORE");
params.add(destName);
return commandExecutor.writeAsync(getRawName(), codec, RedisCommands.SORT_TO, params.toArray());
}
private <T> RFuture<Collection<T>> readSortAsync(String byPattern, List<String> getPatterns, SortOrder order, int offset, int count, boolean alpha) {
List<Object> params = new ArrayList<Object>();
params.add(getRawName());
if (byPattern != null) {
params.add("BY");
params.add(byPattern);
}
if (offset != -1 && count != -1) {
params.add("LIMIT");
}
if (offset != -1) {
params.add(offset);
}
if (count != -1) {
params.add(count);
}
if (getPatterns != null) {
for (String pattern : getPatterns) {
params.add("GET");
params.add(pattern);
}
}
if (alpha) {
params.add("ALPHA");
}
if (order != null) {
params.add(order);
}
return commandExecutor.readAsync(getRawName(), codec, RedisCommands.SORT_LIST, params.toArray());
}
@Override
public RFuture<List<V>> rangeAsync(int toIndex) {
return rangeAsync(0, toIndex);
}
@Override
public RFuture<List<V>> rangeAsync(int fromIndex, int toIndex) {
return commandExecutor.readAsync(getRawName(), codec, LRANGE, getRawName(), fromIndex, toIndex);
}
@Override
public List<V> range(int toIndex) {
return get(rangeAsync(toIndex));
}
@Override
public List<V> range(int fromIndex, int toIndex) {
return get(rangeAsync(fromIndex, toIndex));
}
@Override
public int addListener(ObjectListener listener) {
if (listener instanceof ListAddListener) {
return addListener("__keyevent@*:rpush", (ListAddListener) listener, ListAddListener::onListAdd);
}
if (listener instanceof ListRemoveListener) {
return addListener("__keyevent@*:lrem", (ListRemoveListener) listener, ListRemoveListener::onListRemove);
}
if (listener instanceof ListTrimListener) {
return addListener("__keyevent@*:ltrim", (ListTrimListener) listener, ListTrimListener::onListTrim);
}
if (listener instanceof ListSetListener) {
return addListener("__keyevent@*:lset", (ListSetListener) listener, ListSetListener::onListSet);
}
if (listener instanceof ListInsertListener) {
return addListener("__keyevent@*:linsert", (ListInsertListener) listener, ListInsertListener::onListInsert);
}
return super.addListener(listener);
}
@Override
public RFuture<Integer> addListenerAsync(ObjectListener listener) {
if (listener instanceof ListAddListener) {
return addListenerAsync("__keyevent@*:rpush", (ListAddListener) listener, ListAddListener::onListAdd);
}
if (listener instanceof ListRemoveListener) {
return addListenerAsync("__keyevent@*:lrem", (ListRemoveListener) listener, ListRemoveListener::onListRemove);
}
if (listener instanceof ListTrimListener) {
return addListenerAsync("__keyevent@*:ltrim", (ListTrimListener) listener, ListTrimListener::onListTrim);
}
if (listener instanceof ListSetListener) {
return addListenerAsync("__keyevent@*:lset", (ListSetListener) listener, ListSetListener::onListSet);
}
if (listener instanceof ListInsertListener) {
return addListenerAsync("__keyevent@*:linsert", (ListInsertListener) listener, ListInsertListener::onListInsert);
}
return super.addListenerAsync(listener);
}
@Override
public void removeListener(int listenerId) {
RPatternTopic addTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:rpush");
addTopic.removeListener(listenerId);
RPatternTopic remTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lrem");
remTopic.removeListener(listenerId);
RPatternTopic trimTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:ltrim");
trimTopic.removeListener(listenerId);
RPatternTopic setTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lset");
setTopic.removeListener(listenerId);
RPatternTopic insertTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:linsert");
insertTopic.removeListener(listenerId);
super.removeListener(listenerId);
}
@Override
public RFuture<Void> removeListenerAsync(int listenerId) {
RPatternTopic addTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:rpush");
RFuture<Void> f1 = addTopic.removeListenerAsync(listenerId);
RPatternTopic remTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lrem");
RFuture<Void> f2 = remTopic.removeListenerAsync(listenerId);
RPatternTopic trimTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:ltrim");
RFuture<Void> f3 = trimTopic.removeListenerAsync(listenerId);
RPatternTopic setTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:lset");
RFuture<Void> f4 = setTopic.removeListenerAsync(listenerId);
RPatternTopic insertTopic = new RedissonPatternTopic(StringCodec.INSTANCE, commandExecutor, "__keyevent@*:linsert");
RFuture<Void> f5 = insertTopic.removeListenerAsync(listenerId);
RFuture<Void> f6 = super.removeListenerAsync(listenerId);
CompletableFuture<Void> f = CompletableFuture.allOf(f1.toCompletableFuture(), f2.toCompletableFuture(), f3.toCompletableFuture(),
f4.toCompletableFuture(), f5.toCompletableFuture(), f5.toCompletableFuture(), f6.toCompletableFuture());
return new CompletableFutureWrapper<>(f);
}
@Override
public boolean removeIf(Predicate<? super V> filter) {
throw new UnsupportedOperationException();
}
}

@ -233,13 +233,27 @@ public class RedissonLocalCachedMap<K, V> extends RedissonMap<K, V> implements R
return new CompletableFutureWrapper<>(f);
}
CompletableFuture<V> promise = new CompletableFuture<>();
promise.thenAccept(value -> {
if (storeCacheMiss || value != null) {
cachePut(cacheKey, key, value);
String name = getRawName(key);
RFuture<Boolean> future = containsKeyOperationAsync(name, key);
CompletionStage<Boolean> result = future.thenCompose(res -> {
if (hasNoLoader()) {
if (!res && storeCacheMiss) {
cachePut(cacheKey, key, null);
}
return CompletableFuture.completedFuture(res);
}
if (!res) {
CompletableFuture<V> f = loadValue((K) key, false);
return f.thenApply(value -> {
if (storeCacheMiss || value != null) {
cachePut(cacheKey, key, value);
}
return value != null;
});
}
return CompletableFuture.completedFuture(res);
});
return containsKeyAsync(key, promise);
return new CompletableFutureWrapper<>(result);
}
return new CompletableFutureWrapper<>(cacheValue.getValue() != null);

@ -329,7 +329,4 @@ public class RedissonPriorityDeque<V> extends RedissonPriorityQueue<V> implement
throw new UnsupportedOperationException();
}
public RedissonPriorityDeque<V> reversed() {
throw new UnsupportedOperationException();
}
}

@ -38,7 +38,7 @@ import java.util.function.Supplier;
*
* @param <V> value type
*/
public class RedissonPriorityQueue<V> extends RedissonList<V> implements RPriorityQueue<V> {
public class RedissonPriorityQueue<V> extends BaseRedissonList<V> implements RPriorityQueue<V> {
public static class BinarySearchResult<V> {

@ -33,7 +33,7 @@ import java.util.NoSuchElementException;
*
* @param <V> the type of elements held in this collection
*/
public class RedissonQueue<V> extends RedissonList<V> implements RQueue<V> {
public class RedissonQueue<V> extends BaseRedissonList<V> implements RQueue<V> {
public RedissonQueue(CommandAsyncExecutor commandExecutor, String name, RedissonClient redisson) {
super(commandExecutor, name, redisson);

@ -22,13 +22,17 @@ import org.redisson.api.search.aggregate.*;
import org.redisson.api.search.index.*;
import org.redisson.api.search.query.*;
import org.redisson.client.codec.Codec;
import org.redisson.client.codec.DoubleCodec;
import org.redisson.client.codec.LongCodec;
import org.redisson.client.codec.StringCodec;
import org.redisson.client.protocol.RedisCommand;
import org.redisson.client.protocol.RedisCommands;
import org.redisson.client.protocol.RedisStrictCommand;
import org.redisson.client.protocol.convertor.EmptyMapConvertor;
import org.redisson.client.protocol.decoder.*;
import org.redisson.codec.CompositeCodec;
import org.redisson.command.CommandAsyncExecutor;
import org.redisson.config.Protocol;
import java.math.BigDecimal;
import java.util.ArrayList;
@ -151,7 +155,7 @@ public class RedissonSearch implements RSearch {
}
args.add("VECTOR");
args.add("HNSW");
args.add(params.getCount());
args.add(params.getCount()*2);
args.add("TYPE");
args.add(params.getType());
args.add("DIM");
@ -473,14 +477,27 @@ public class RedissonSearch implements RSearch {
args.add(options.getDialect());
}
RedisStrictCommand<SearchResult> command = new RedisStrictCommand<>("FT.SEARCH",
new ListMultiDecoder2(new SearchResultDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec)),
new ObjectListReplayDecoder()));
RedisStrictCommand<SearchResult> command;
if (isResp3()) {
command = new RedisStrictCommand<>("FT.SEARCH",
new ListMultiDecoder2(new SearchResultDecoderV2(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
} else {
command = new RedisStrictCommand<>("FT.SEARCH",
new ListMultiDecoder2(new SearchResultDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec)),
new ObjectListReplayDecoder()));
}
return commandExecutor.writeAsync(indexName, StringCodec.INSTANCE, command, args.toArray());
}
private boolean isResp3() {
return commandExecutor.getServiceManager().getCfg().getProtocol() == Protocol.RESP3;
}
private String value(double score, boolean exclusive) {
StringBuilder element = new StringBuilder();
if (Double.isInfinite(score)) {
@ -593,19 +610,35 @@ public class RedissonSearch implements RSearch {
}
RedisStrictCommand<AggregationResult> command;
if (options.isWithCursor()) {
command = new RedisStrictCommand<>("FT.AGGREGATE",
new ListMultiDecoder2(new AggregationCursorResultDecoder(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
if (isResp3()) {
if (options.isWithCursor()) {
command = new RedisStrictCommand<>("FT.AGGREGATE",
new ListMultiDecoder2(new AggregationCursorResultDecoderV2(),
new ObjectListReplayDecoder(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
} else {
command = new RedisStrictCommand<>("FT.AGGREGATE",
new ListMultiDecoder2(new AggregationResultDecoderV2(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
}
} else {
command = new RedisStrictCommand<>("FT.AGGREGATE",
new ListMultiDecoder2(new AggregationResultDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec)),
new ObjectListReplayDecoder()));
if (options.isWithCursor()) {
command = new RedisStrictCommand<>("FT.AGGREGATE",
new ListMultiDecoder2(new AggregationCursorResultDecoder(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
} else {
command = new RedisStrictCommand<>("FT.AGGREGATE",
new ListMultiDecoder2(new AggregationResultDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec)),
new ObjectListReplayDecoder()));
}
}
return commandExecutor.writeAsync(indexName, StringCodec.INSTANCE, command, args.toArray());
}
@ -703,10 +736,20 @@ public class RedissonSearch implements RSearch {
@Override
public RFuture<AggregationResult> readCursorAsync(String indexName, long cursorId) {
RedisStrictCommand command = new RedisStrictCommand<>("FT.CURSOR", "READ",
new ListMultiDecoder2(new AggregationCursorResultDecoder(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
RedisStrictCommand command;
if (isResp3()) {
command = new RedisStrictCommand<>("FT.CURSOR", "READ",
new ListMultiDecoder2(new AggregationCursorResultDecoderV2(),
new ObjectListReplayDecoder(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
} else {
command = new RedisStrictCommand<>("FT.CURSOR", "READ",
new ListMultiDecoder2(new AggregationCursorResultDecoder(),
new ObjectListReplayDecoder(),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, codec))));
}
return commandExecutor.writeAsync(indexName, StringCodec.INSTANCE, command, indexName, cursorId);
}
@ -824,7 +867,17 @@ public class RedissonSearch implements RSearch {
args.add(options.getDialect());
}
return commandExecutor.readAsync(indexName, StringCodec.INSTANCE, RedisCommands.FT_SPELLCHECK, args.toArray());
RedisCommand<Map<String, Map<String, Object>>> command = RedisCommands.FT_SPELLCHECK;
if (isResp3()) {
command = new RedisCommand<>("FT.SPELLCHECK",
new ListMultiDecoder2(
new ListObjectDecoder(1),
new ObjectMapReplayDecoder(),
new ListFirstObjectDecoder(new EmptyMapConvertor()),
new ObjectMapReplayDecoder(new CompositeCodec(StringCodec.INSTANCE, DoubleCodec.INSTANCE))));
}
return commandExecutor.readAsync(indexName, StringCodec.INSTANCE, command, args.toArray());
}
@Override

@ -51,7 +51,6 @@ import java.net.InetSocketAddress;
import java.net.SocketOption;
import java.net.UnknownHostException;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
/**
@ -80,11 +79,6 @@ public final class RedisClient {
private boolean hasOwnResolver;
private volatile boolean shutdown;
private final AtomicLong firstFailTime = new AtomicLong(0);
private Runnable connectedListener;
private Runnable disconnectedListener;
public static RedisClient create(RedisClientConfig config) {
return new RedisClient(config);
}

@ -20,10 +20,7 @@ import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.resolver.AddressResolverGroup;
import io.netty.util.Timer;
import org.redisson.config.CommandMapper;
import org.redisson.config.CredentialsResolver;
import org.redisson.config.DefaultCommandMapper;
import org.redisson.config.SslProvider;
import org.redisson.config.*;
import org.redisson.misc.RedisURI;
import javax.net.ssl.KeyManagerFactory;
@ -85,6 +82,8 @@ public class RedisClientConfig {
private FailedNodeDetector failedNodeDetector = new FailedConnectionDetector();
private Protocol protocol = Protocol.RESP2;
public RedisClientConfig() {
}
@ -129,6 +128,7 @@ public class RedisClientConfig {
this.tcpKeepAliveIdle = config.tcpKeepAliveIdle;
this.tcpKeepAliveInterval = config.tcpKeepAliveInterval;
this.tcpUserTimeout = config.tcpUserTimeout;
this.protocol = config.protocol;
}
public NettyHook getNettyHook() {
@ -457,4 +457,13 @@ public class RedisClientConfig {
this.failedNodeDetector = failedNodeDetector;
return this;
}
public Protocol getProtocol() {
return protocol;
}
public RedisClientConfig setProtocol(Protocol protocol) {
this.protocol = protocol;
return this;
}
}

@ -19,6 +19,7 @@ import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import org.redisson.client.*;
import org.redisson.client.protocol.RedisCommands;
import org.redisson.config.Protocol;
import java.net.InetSocketAddress;
import java.util.ArrayList;
@ -79,8 +80,10 @@ public abstract class BaseConnectionHandler<C extends RedisConnection> extends C
});
futures.add(f.toCompletableFuture());
// CompletionStage<Object> f1 = connection.async(RedisCommands.HELLO, "3");
// futures.add(f1.toCompletableFuture());
if (redisClient.getConfig().getProtocol() == Protocol.RESP3) {
CompletionStage<Object> f1 = connection.async(RedisCommands.HELLO, "3");
futures.add(f1.toCompletableFuture());
}
if (config.getDatabase() != 0) {
CompletionStage<Object> future = connection.async(RedisCommands.SELECT, config.getDatabase());

@ -164,7 +164,11 @@ public class CommandDecoder extends ReplayingDecoder<State> {
protected void skipDecode(ByteBuf in) throws IOException{
int code = in.readByte();
if (code == '+') {
if (code == '_') {
in.skipBytes(2);
} else if (code == ',') {
skipString(in);
} else if (code == '+') {
skipString(in);
} else if (code == '-') {
skipString(in);
@ -172,7 +176,14 @@ public class CommandDecoder extends ReplayingDecoder<State> {
skipString(in);
} else if (code == '$') {
skipBytes(in);
} else if (code == '*') {
} else if (code == '=') {
skipBytes(in);
} else if (code == '%') {
long size = readLong(in);
for (int i = 0; i < size * 2; i++) {
skipDecode(in);
}
} else if (code == '*' || code == '>' || code == '~') {
long size = readLong(in);
for (int i = 0; i < size; i++) {
skipDecode(in);
@ -335,9 +346,21 @@ public class CommandDecoder extends ReplayingDecoder<State> {
protected void decode(ByteBuf in, CommandData<Object, Object> data, List<Object> parts, Channel channel, boolean skipConvertor, List<CommandData<?, ?>> commandsData) throws IOException {
int code = in.readByte();
if (code == '+') {
if (code == '_') {
readCRLF(in);
Object result = null;
handleResult(data, parts, result, false);
} else if (code == '+') {
String result = readString(in);
handleResult(data, parts, result, skipConvertor);
} else if (code == ',') {
String str = readString(in);
Double result = Double.NaN;
if (!"nan".equals(str)) {
result = Double.valueOf(str);
}
handleResult(data, parts, result, skipConvertor);
} else if (code == '-') {
String error = readString(in);
@ -386,6 +409,15 @@ public class CommandDecoder extends ReplayingDecoder<State> {
} else if (code == ':') {
Long result = readLong(in);
handleResult(data, parts, result, false);
} else if (code == '=') {
ByteBuf buf = readBytes(in);
Object result = null;
if (buf != null) {
buf.skipBytes(3);
Decoder<Object> decoder = selectDecoder(data, parts);
result = decoder.decode(buf, state());
}
handleResult(data, parts, result, false);
} else if (code == '$') {
ByteBuf buf = readBytes(in);
Object result = null;
@ -394,7 +426,7 @@ public class CommandDecoder extends ReplayingDecoder<State> {
result = decoder.decode(buf, state());
}
handleResult(data, parts, result, false);
} else if (code == '*') {
} else if (code == '*' || code == '>' || code == '~') {
long size = readLong(in);
List<Object> respParts = new ArrayList<Object>(Math.max((int) size, 0));
@ -403,7 +435,16 @@ public class CommandDecoder extends ReplayingDecoder<State> {
decodeList(in, data, parts, channel, size, respParts, skipConvertor, commandsData);
state().decLevel();
} else if (code == '%') {
long size = readLong(in) * 2;
List<Object> respParts = new ArrayList<Object>(Math.max((int) size, 0));
state().incLevel();
decodeList(in, data, parts, channel, size, respParts, skipConvertor, commandsData);
state().decLevel();
} else {
String dataStr = in.toString(0, in.writerIndex(), CharsetUtil.UTF_8);
throw new IllegalStateException("Can't decode replay: " + dataStr);
@ -420,7 +461,7 @@ public class CommandDecoder extends ReplayingDecoder<State> {
in.skipBytes(len + 2);
return result;
}
@SuppressWarnings("unchecked")
private void decodeList(ByteBuf in, CommandData<Object, Object> data, List<Object> parts,
Channel channel, long size, List<Object> respParts, boolean skipConvertor, List<CommandData<?, ?>> commandsData)
@ -514,6 +555,10 @@ public class CommandDecoder extends ReplayingDecoder<State> {
return buffer;
}
private void readCRLF(ByteBuf is) {
is.skipBytes(2);
}
private long readLong(ByteBuf is) throws IOException {
long size = 0;
int sign = 1;

@ -204,7 +204,7 @@ public class CommandPubSubDecoder extends CommandDecoder {
});
}
} else {
if (data != null && data.getCommand().getName().equals("PING")) {
if (data != null) {
super.decodeResult(data, parts, channel, result);
}
}

@ -282,7 +282,7 @@ public interface RedisCommands {
RedisCommand<List<Object>> BLPOP = new RedisCommand<List<Object>>("BLPOP", new ObjectListReplayDecoder<Object>());
RedisCommand<List<Object>> BRPOP = new RedisCommand<List<Object>>("BRPOP", new ObjectListReplayDecoder<Object>());
RedisCommand<Map<String, Map<Object, Double>>> BZMPOP = new RedisCommand<>("BZMPOP", ZMPOP.getReplayMultiDecoder());
RedisCommand<List<Object>> BZMPOP_SINGLE_LIST = new RedisCommand<>("BZMPOP", ZMPOP_VALUES.getReplayMultiDecoder());
RedisCommand<List<Object>> BZMPOP_SINGLE_LIST = new RedisCommand("BZMPOP", ZMPOP_VALUES.getReplayMultiDecoder(), new EmptyListConvertor());
RedisCommand<Object> BLPOP_VALUE = new RedisCommand<Object>("BLPOP", new ListObjectDecoder<Object>(1));
RedisCommand<Object> BLMOVE = new RedisCommand<Object>("BLMOVE");
RedisCommand<Object> BRPOP_VALUE = new RedisCommand<Object>("BRPOP", new ListObjectDecoder<Object>(1));

@ -0,0 +1,60 @@
/**
* Copyright (c) 2013-2022 Nikita Koksharov
*
* 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.redisson.client.protocol.convertor;
import java.math.BigDecimal;
/**
*
* @author Nikita Koksharov
*
*/
public class LongNumberConvertor implements Convertor<Object> {
private Class<?> resultClass;
public LongNumberConvertor(Class<?> resultClass) {
super();
this.resultClass = resultClass;
}
@Override
public Object convert(Object result) {
if (result instanceof Long) {
Long res = (Long) result;
if (resultClass.isAssignableFrom(Long.class)) {
return res;
}
if (resultClass.isAssignableFrom(Integer.class)) {
return res.intValue();
}
if (resultClass.isAssignableFrom(BigDecimal.class)) {
return new BigDecimal(res);
}
}
if (result instanceof Double) {
Double res = (Double) result;
if (resultClass.isAssignableFrom(Float.class)) {
return ((Double) result).floatValue();
}
if (resultClass.isAssignableFrom(Double.class)) {
return res;
}
}
throw new IllegalStateException("Wrong value type!");
}
}

@ -0,0 +1,55 @@
/**
* Copyright (c) 2013-2022 Nikita Koksharov
*
* 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.redisson.client.protocol.decoder;
import org.redisson.api.search.aggregate.AggregationResult;
import org.redisson.client.handler.State;
import java.util.*;
/**
*
* @author Nikita Koksharov
*
*/
public class AggregationCursorResultDecoderV2 implements MultiDecoder<Object> {
@Override
public Object decode(List<Object> parts, State state) {
if (parts.isEmpty()) {
return new AggregationResult(0, Collections.emptyList(), -1);
}
List<Object> attrs = (List<Object>) parts.get(0);
Map<String, Object> m = new HashMap<>();
for (int i = 0; i < attrs.size(); i++) {
if (i % 2 != 0) {
m.put(attrs.get(i-1).toString(), attrs.get(i));
}
}
List<Map<String, Object>> docs = new ArrayList<>();
List<Map<String, Object>> results = (List<Map<String, Object>>) m.get("results");
for (Map<String, Object> result : results) {
Map<String, Object> map = (Map<String, Object>) result.get("extra_attributes");
docs.add(map);
}
Long total = (Long) m.get("total_results");
long cursorId = (long) parts.get(1);
return new AggregationResult(total, docs, cursorId);
}
}

@ -0,0 +1,56 @@
/**
* Copyright (c) 2013-2022 Nikita Koksharov
*
* 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.redisson.client.protocol.decoder;
import org.redisson.api.search.aggregate.AggregationResult;
import org.redisson.client.handler.State;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
*
* @author Nikita Koksharov
*
*/
public class AggregationResultDecoderV2 implements MultiDecoder<Object> {
@Override
public Object decode(List<Object> parts, State state) {
if (parts.isEmpty()) {
return null;
}
Map<String, Object> m = new HashMap<>();
for (int i = 0; i < parts.size(); i++) {
if (i % 2 != 0) {
m.put(parts.get(i-1).toString(), parts.get(i));
}
}
List<Map<String, Object>> docs = new ArrayList<>();
List<Map<String, Object>> results = (List<Map<String, Object>>) m.get("results");
for (Map<String, Object> result : results) {
Map<String, Object> attrs = (Map<String, Object>) result.get("extra_attributes");
docs.add(attrs);
}
Long total = (Long) m.get("total_results");
return new AggregationResult(total, docs);
}
}

@ -79,6 +79,10 @@ public class IndexInfoDecoder implements MultiDecoder<Object> {
if (result.get(prop).toString().contains("nan")) {
return 0L;
}
if (result.get(prop) instanceof Double) {
Double d = (Double) result.get(prop);
return d.longValue();
}
return Long.valueOf(result.get(prop).toString());
}
}

@ -18,6 +18,7 @@ package org.redisson.client.protocol.decoder;
import org.redisson.client.codec.Codec;
import org.redisson.client.handler.State;
import org.redisson.client.protocol.Decoder;
import org.redisson.client.protocol.ScoredEntry;
import org.redisson.client.protocol.convertor.Convertor;
import java.util.List;
@ -54,7 +55,7 @@ public class ListFirstObjectDecoder implements MultiDecoder<Object> {
@Override
public Object decode(List<Object> parts, State state) {
if (inner != null) {
if (inner != null && !parts.isEmpty() && !(parts.get(0) instanceof ScoredEntry)) {
parts = (List) inner.decode(parts, state);
}
if (!parts.isEmpty()) {

@ -43,7 +43,7 @@ public class ObjectFirstScoreReplayDecoder implements MultiDecoder<Double> {
if (parts.isEmpty()) {
return null;
}
return (Double) parts.get(1);
return (Double) parts.get(parts.size()-1);
}
}

@ -22,6 +22,7 @@ import org.redisson.client.protocol.Decoder;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
/**
*
@ -57,6 +58,13 @@ public class ObjectMapReplayDecoder<K, V> implements MultiDecoder<Map<K, V>> {
@Override
public Map<K, V> decode(List<Object> parts, State state) {
if (!parts.isEmpty() && parts.get(0) instanceof Map) {
return ((List<Map<K, V>>) (Object) parts)
.stream()
.flatMap(v -> v.entrySet().stream())
.collect(Collectors.toMap(v -> v.getKey(), v -> v.getValue()));
}
Map<K, V> result = MultiDecoder.newLinkedHashMap(parts.size()/2);
for (int i = 0; i < parts.size(); i++) {
if (i % 2 != 0) {

@ -20,6 +20,10 @@ import org.redisson.client.codec.DoubleCodec;
import org.redisson.client.handler.State;
import org.redisson.client.protocol.Decoder;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
/**
*
* @author Nikita Koksharov
@ -35,4 +39,15 @@ public class ScoredSortedSetRandomMapDecoder extends ObjectMapReplayDecoder<Obje
return DoubleCodec.INSTANCE.getValueDecoder();
}
@Override
public Map<Object, Object> decode(List<Object> parts, State state) {
if (!parts.isEmpty() && parts.get(0) instanceof Map) {
return ((List<Map<Object, Object>>) (Object) parts)
.stream()
.flatMap(v -> v.entrySet().stream())
.collect(Collectors.toMap(v -> v.getKey(), v -> v.getValue()));
}
return super.decode(parts, state);
}
}

@ -17,6 +17,7 @@ package org.redisson.client.protocol.decoder;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Collectors;
import org.redisson.client.codec.Codec;
import org.redisson.client.codec.DoubleCodec;
@ -42,6 +43,9 @@ public class ScoredSortedSetReplayDecoder<T> implements MultiDecoder<List<Scored
@Override
public List<ScoredEntry<T>> decode(List<Object> parts, State state) {
if (!parts.isEmpty() && parts.get(0) instanceof List) {
return ((List<List<ScoredEntry<T>>>) (Object) parts).stream().flatMap(v -> v.stream()).collect(Collectors.toList());
}
List<ScoredEntry<T>> result = new ArrayList<>();
for (int i = 0; i < parts.size(); i += 2) {
result.add(new ScoredEntry<T>(((Number) parts.get(i+1)).doubleValue(), (T) parts.get(i)));

@ -0,0 +1,55 @@
/**
* Copyright (c) 2013-2022 Nikita Koksharov
*
* 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.redisson.client.protocol.decoder;
import org.redisson.api.search.query.Document;
import org.redisson.api.search.query.SearchResult;
import org.redisson.client.handler.State;
import java.util.*;
/**
*
* @author Nikita Koksharov
*
*/
public class SearchResultDecoderV2 implements MultiDecoder<Object> {
@Override
public Object decode(List<Object> parts, State state) {
if (parts.isEmpty()) {
return new SearchResult(0, Collections.emptyList());
}
Map<String, Object> m = new HashMap<>();
for (int i = 0; i < parts.size(); i++) {
if (i % 2 != 0) {
m.put(parts.get(i-1).toString(), parts.get(i));
}
}
List<Document> docs = new ArrayList<>();
List<Map<String, Object>> results = (List<Map<String, Object>>) m.get("results");
for (Map<String, Object> result : results) {
String id = (String) result.get("id");
Map<String, Object> attrs = (Map<String, Object>) result.get("extra_attributes");
docs.add(new Document(id, attrs));
}
Long total = (Long) m.get("total_results");
return new SearchResult(total, docs);
}
}

@ -51,7 +51,7 @@ import static com.esotericsoftware.kryo.util.Util.className;
*/
public class Kryo5Codec extends BaseCodec {
private static class SimpleInstantiatorStrategy implements org.objenesis.strategy.InstantiatorStrategy {
private static final class SimpleInstantiatorStrategy implements org.objenesis.strategy.InstantiatorStrategy {
private final StdInstantiatorStrategy ss = new StdInstantiatorStrategy();

@ -183,7 +183,7 @@ public class ProtobufCodec extends BaseCodec {
};
}
private static class ProtostuffUtils {
private static final class ProtostuffUtils {
@SuppressWarnings("unchecked")
public static <T> byte[] serialize(T obj) {

@ -140,6 +140,8 @@ public interface CommandAsyncExecutor {
<T, R> RFuture<R> evalWriteBatchedAsync(Codec codec, RedisCommand<T> command, String script, List<Object> keys, SlotCallback<T, R> callback);
<T, R> RFuture<R> evalReadBatchedAsync(Codec codec, RedisCommand<T> command, String script, List<Object> keys, SlotCallback<T, R> callback);
boolean isEvalShaROSupported();
void setEvalShaROSupported(boolean value);

@ -583,10 +583,15 @@ public class CommandAsyncService implements CommandAsyncExecutor {
@Override
public <T, R> RFuture<R> evalWriteBatchedAsync(Codec codec, RedisCommand<T> command, String script, List<Object> keys, SlotCallback<T, R> callback) {
return evalWriteBatchedAsync(false, codec, command, script, keys, callback);
return evalBatchedAsync(false, codec, command, script, keys, callback);
}
private <T, R> RFuture<R> evalWriteBatchedAsync(boolean readOnly, Codec codec, RedisCommand<T> command, String script, List<Object> keys, SlotCallback<T, R> callback) {
@Override
public <T, R> RFuture<R> evalReadBatchedAsync(Codec codec, RedisCommand<T> command, String script, List<Object> keys, SlotCallback<T, R> callback) {
return evalBatchedAsync(true, codec, command, script, keys, callback);
}
private <T, R> RFuture<R> evalBatchedAsync(boolean readOnly, Codec codec, RedisCommand<T> command, String script, List<Object> keys, SlotCallback<T, R> callback) {
if (!connectionManager.isClusterMode()) {
Object[] keysArray = callback.createKeys(keys);
Object[] paramsArray = callback.createParams(Collections.emptyList());

@ -378,7 +378,11 @@ public class RedisExecutor<V, R> {
}
}
} else {
popTimeout = Long.valueOf(params[params.length - 1].toString()) * 1000;
if (RedisCommands.BZMPOP.getName().equals(command.getName())) {
popTimeout = Long.valueOf(params[0].toString()) * 1000;
} else {
popTimeout = Long.valueOf(params[params.length - 1].toString()) * 1000;
}
}
handleBlockingOperations(attemptPromise, connection, popTimeout);

@ -96,6 +96,8 @@ public class Config {
private boolean lazyInitialization;
private Protocol protocol = Protocol.RESP2;
public Config() {
}
@ -127,6 +129,7 @@ public class Config {
setAddressResolverGroupFactory(oldConf.getAddressResolverGroupFactory());
setReliableTopicWatchdogTimeout(oldConf.getReliableTopicWatchdogTimeout());
setLazyInitialization(oldConf.isLazyInitialization());
setProtocol(oldConf.getProtocol());
if (oldConf.getSingleServerConfig() != null) {
setSingleServerConfig(new SingleServerConfig(oldConf.getSingleServerConfig()));
@ -879,10 +882,27 @@ public class Config {
*
* @param lazyInitialization <code>true</code> connects to Redis only when first Redis call is made,
* <code>false</code> connects to Redis during Redisson instance creation.
* @return
* @return config
*/
public Config setLazyInitialization(boolean lazyInitialization) {
this.lazyInitialization = lazyInitialization;
return this;
}
public Protocol getProtocol() {
return protocol;
}
/**
* Defines Redis protocol version.
* <p>
* Default value is <code>RESP2</code>
*
* @param protocol Redis protocol version
* @return config
*/
public Config setProtocol(Protocol protocol) {
this.protocol = protocol;
return this;
}
}

@ -0,0 +1,30 @@
/**
* Copyright (c) 2013-2022 Nikita Koksharov
*
* 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.redisson.config;
/**
* Redis protocol version
*
* @author Nikita Koksharov
*
*/
public enum Protocol {
RESP2,
RESP3
}

@ -361,6 +361,7 @@ public class MasterSlaveConnectionManager implements ConnectionManager {
.setPassword(config.getPassword())
.setNettyHook(serviceManager.getCfg().getNettyHook())
.setFailedNodeDetector(config.getFailedSlaveNodeDetector())
.setProtocol(serviceManager.getCfg().getProtocol())
.setCommandMapper(config.getCommandMapper())
.setCredentialsResolver(config.getCredentialsResolver())
.setConnectedListener(addr -> {

@ -15,14 +15,13 @@
*/
package org.redisson.reactive;
import java.util.concurrent.Callable;
import org.redisson.BaseRedissonList;
import org.redisson.api.RBlockingQueue;
import org.redisson.api.RFuture;
import org.redisson.api.RListAsync;
import reactor.core.publisher.Flux;
import java.util.concurrent.Callable;
/**
*
* @author Nikita Koksharov
@ -34,7 +33,7 @@ public class RedissonBlockingQueueReactive<V> extends RedissonListReactive<V> {
private final RBlockingQueue<V> queue;
public RedissonBlockingQueueReactive(RBlockingQueue<V> queue) {
super((RListAsync<V>) queue);
super((BaseRedissonList<V>) queue);
this.queue = queue;
}

@ -15,18 +15,17 @@
*/
package org.redisson.reactive;
import java.util.function.Consumer;
import java.util.function.LongConsumer;
import org.reactivestreams.Publisher;
import org.redisson.BaseRedissonList;
import org.redisson.RedissonList;
import org.redisson.api.RFuture;
import org.redisson.api.RListAsync;
import org.redisson.client.codec.Codec;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import java.util.function.Consumer;
import java.util.function.LongConsumer;
/**
* Distributed and concurrent implementation of {@link java.util.List}
*
@ -36,9 +35,9 @@ import reactor.core.publisher.FluxSink;
*/
public class RedissonListReactive<V> {
private final RListAsync<V> instance;
private final BaseRedissonList<V> instance;
public RedissonListReactive(RListAsync<V> instance) {
public RedissonListReactive(BaseRedissonList<V> instance) {
this.instance = instance;
}

@ -15,10 +15,9 @@
*/
package org.redisson.rx;
import org.redisson.api.RBlockingQueueAsync;
import org.redisson.api.RListAsync;
import io.reactivex.rxjava3.core.Flowable;
import org.redisson.BaseRedissonList;
import org.redisson.api.RBlockingQueueAsync;
/**
*
@ -31,7 +30,7 @@ public class RedissonBlockingQueueRx<V> extends RedissonListRx<V> {
private final RBlockingQueueAsync<V> queue;
public RedissonBlockingQueueRx(RBlockingQueueAsync<V> queue) {
super((RListAsync<V>) queue);
super((BaseRedissonList<V>) queue);
this.queue = queue;
}

@ -15,13 +15,12 @@
*/
package org.redisson.rx;
import org.reactivestreams.Publisher;
import org.redisson.api.RFuture;
import org.redisson.api.RListAsync;
import io.reactivex.rxjava3.core.Single;
import io.reactivex.rxjava3.functions.LongConsumer;
import io.reactivex.rxjava3.processors.ReplayProcessor;
import org.reactivestreams.Publisher;
import org.redisson.BaseRedissonList;
import org.redisson.api.RFuture;
/**
* Distributed and concurrent implementation of {@link java.util.List}
@ -32,9 +31,9 @@ import io.reactivex.rxjava3.processors.ReplayProcessor;
*/
public class RedissonListRx<V> {
private final RListAsync<V> instance;
private final BaseRedissonList<V> instance;
public RedissonListRx(RListAsync<V> instance) {
public RedissonListRx(BaseRedissonList<V> instance) {
this.instance = instance;
}

@ -155,8 +155,6 @@ public abstract class BaseMapTest extends BaseTest {
@Test
public void testRandomEntries() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RMap<Integer, Integer> map = getMap("map");
Map<Integer, Integer> e1 = map.randomEntries(1);
assertThat(e1).isEmpty();

@ -0,0 +1,45 @@
package org.redisson;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.redisson.api.RedissonClient;
import org.redisson.config.Config;
import org.testcontainers.containers.GenericContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
@Testcontainers
public class DockerRedisStackTest {
@Container
private static final GenericContainer<?> REDIS =
new GenericContainer<>("redis/redis-stack")
.withExposedPorts(6379);
protected static RedissonClient redisson;
@BeforeAll
public static void beforeAll() {
Config config = createConfig();
redisson = Redisson.create(config);
}
protected static Config createConfig() {
Config config = new Config();
config.useSingleServer()
.setAddress("redis://127.0.0.1:" + REDIS.getFirstMappedPort());
return config;
}
@BeforeEach
public void beforeEach() {
redisson.getKeys().flushall();
}
@AfterAll
public static void afterAll() {
redisson.shutdown();
}
}

@ -0,0 +1,85 @@
package org.redisson;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.redisson.api.*;
import org.redisson.config.Config;
import org.redisson.config.Protocol;
import org.redisson.misc.RedisURI;
import org.testcontainers.containers.GenericContainer;
import org.testcontainers.containers.startupcheck.MinimumDurationRunningStartupCheckStrategy;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import java.time.Duration;
import java.util.function.Consumer;
import static org.assertj.core.api.Assertions.assertThat;
@Testcontainers
public class RedisDockerTest {
@Container
private static final GenericContainer<?> REDIS =
new GenericContainer<>("redis:7.2")
.withExposedPorts(6379);
protected static RedissonClient redisson;
@BeforeAll
public static void beforeAll() {
Config config = createConfig();
redisson = Redisson.create(config);
}
protected static Config createConfig() {
Config config = new Config();
config.setProtocol(Protocol.RESP3);
config.useSingleServer()
.setAddress("redis://127.0.0.1:" + REDIS.getFirstMappedPort());
return config;
}
protected void testInCluster(Consumer<RedissonClient> redissonCallback) {
GenericContainer<?> redisClusterContainer =
new GenericContainer<>("vishnunair/docker-redis-cluster")
.withExposedPorts(6379, 6380, 6381, 6382, 6383, 6384)
.withStartupCheckStrategy(new MinimumDurationRunningStartupCheckStrategy(Duration.ofSeconds(7)));
redisClusterContainer.start();
Config config = new Config();
config.setProtocol(Protocol.RESP3);
config.useClusterServers()
.setNatMapper(new NatMapper() {
@Override
public RedisURI map(RedisURI uri) {
if (redisClusterContainer.getMappedPort(uri.getPort()) == null) {
return uri;
}
return new RedisURI(uri.getScheme(), redisClusterContainer.getHost(), redisClusterContainer.getMappedPort(uri.getPort()));
}
})
.addNodeAddress("redis://127.0.0.1:" + redisClusterContainer.getFirstMappedPort());
RedissonClient redisson = Redisson.create(config);
try {
redissonCallback.accept(redisson);
} finally {
redisson.shutdown();
redisClusterContainer.stop();
}
}
@BeforeEach
public void beforeEach() {
redisson.getKeys().flushall();
}
@AfterAll
public static void afterAll() {
redisson.shutdown();
}
}

@ -201,7 +201,7 @@ public class RedisRunner {
private boolean randomDir = false;
private ArrayList<String> bindAddr = new ArrayList<>();
private int port = 6379;
private int retryCount = Integer.MAX_VALUE;
private int retryCount = 10;
private boolean randomPort = false;
private String sentinelFile;
private String clusterFile;
@ -294,7 +294,16 @@ public class RedisRunner {
throw new FailedToStartRedisException();
}
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
rp.stop();
if (RedissonRuntimeEnvironment.isWindows
&& RedissonRuntimeEnvironment.redisBinaryPath.contains("cmd")) {
try {
Runtime.getRuntime().exec("C:\\redis\\redis-server-stop.cmd");
} catch (IOException e) {
throw new RuntimeException(e);
}
} else {
rp.stop();
}
}));
return rp;
}
@ -927,31 +936,32 @@ public class RedisRunner {
}
public int stop() {
if (runner.isNosave() && !runner.isRandomDir()) {
RedisClient c = createDefaultRedisClientInstance();
RedisConnection connection = c.connect();
if (runner.isNosave()) {
RedisClientConfig config = new RedisClientConfig();
config.setConnectTimeout(1000);
config.setAddress(runner.getInitialBindAddr(), runner.getPort());
RedisClient c = RedisClient.create(config);
RedisConnection connection = null;
try {
connection.async(new RedisStrictCommand<Void>("SHUTDOWN", "NOSAVE", new VoidReplayConvertor()))
.get(3, TimeUnit.SECONDS);
} catch (InterruptedException interruptedException) {
//shutdown via command failed, lets wait and kill it later.
} catch (ExecutionException | TimeoutException e) {
connection = c.connect();
connection.async(new RedisStrictCommand<Void>("SHUTDOWN", "NOSAVE", new VoidReplayConvertor()));
} catch (Exception e) {
// skip
}
c.shutdown();
connection.closeAsync().syncUninterruptibly();
}
Process p = redisProcess;
p.destroy();
boolean normalTermination = false;
try {
normalTermination = p.waitFor(5, TimeUnit.SECONDS);
} catch (InterruptedException ex) {
//OK lets hurry up by force kill;
}
if (!normalTermination) {
p = p.destroyForcibly();
}
// boolean normalTermination = false;
// try {
// normalTermination = p.waitFor(5, TimeUnit.SECONDS);
// } catch (InterruptedException ex) {
// //OK lets hurry up by force kill;
// }
// if (!normalTermination) {
// p = p.destroyForcibly();
// }
cleanup();
int exitCode = p.exitValue();
return exitCode == 1 && RedissonRuntimeEnvironment.isWindows ? 0 : exitCode;

@ -419,6 +419,7 @@ public class RedissonBatchTest extends BaseTest {
@ParameterizedTest
@MethodSource("data")
@Timeout(20)
public void testSyncSlaves(BatchOptions batchOptions) throws FailedToStartRedisException, IOException, InterruptedException {
RedisRunner master1 = new RedisRunner().randomPort().randomDir().nosave();
RedisRunner master2 = new RedisRunner().randomPort().randomDir().nosave();

@ -15,12 +15,10 @@ import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
public class RedissonBlockingDequeTest extends BaseTest {
public class RedissonBlockingDequeTest extends RedisDockerTest {
@Test
public void testMove() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RBlockingDeque<Integer> deque1 = redisson.getBlockingDeque("deque1");
RBlockingDeque<Integer> deque2 = redisson.getBlockingDeque("deque2");

@ -44,7 +44,7 @@ public class RedissonBlockingQueueReactiveTest extends BaseReactiveTest {
.repeat()
.subscribe();
Awaitility.await().atMost(Duration.ofSeconds(1)).untilAsserted(() -> {
Awaitility.await().atMost(Duration.ofSeconds(2)).untilAsserted(() -> {
assertThat(counter.get()).isEqualTo(100);
});
}

@ -184,7 +184,6 @@ public class RedissonBucketTest extends BaseTest {
@Test
public void testSizeInMemory() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("4.0.0") > 0);
RBucket<Integer> al = redisson.getBucket("test");
al.set(1234);
assertThat(al.sizeInMemory()).isEqualTo(51);

@ -1,28 +1,15 @@
package org.redisson;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.redisson.api.*;
import org.redisson.client.codec.LongCodec;
import org.redisson.client.codec.StringCodec;
import org.redisson.config.Config;
import org.redisson.connection.balancer.RandomLoadBalancer;
import org.redisson.misc.RedisURI;
import org.testcontainers.containers.GenericContainer;
import org.testcontainers.containers.startupcheck.MinimumDurationRunningStartupCheckStrategy;
import java.time.Duration;
import java.util.*;
import static org.assertj.core.api.Assertions.assertThat;
public class RedissonFunctionTest extends BaseTest {
@BeforeAll
public static void check() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
}
public class RedissonFunctionTest extends RedisDockerTest {
@Test
public void testEmpty() {
@ -54,58 +41,41 @@ public class RedissonFunctionTest extends BaseTest {
}
@Test
public void testCluster() throws InterruptedException {
GenericContainer<?> redisClusterContainer =
new GenericContainer<>("vishnunair/docker-redis-cluster")
.withExposedPorts(6379, 6380, 6381, 6382, 6383, 6384)
.withStartupCheckStrategy(new MinimumDurationRunningStartupCheckStrategy(Duration.ofSeconds(7)));
redisClusterContainer.start();
Config config = new Config();
config.useClusterServers()
.setNatMapper(new NatMapper() {
@Override
public RedisURI map(RedisURI uri) {
if (redisClusterContainer.getMappedPort(uri.getPort()) == null) {
return uri;
}
return new RedisURI(uri.getScheme(), redisClusterContainer.getHost(), redisClusterContainer.getMappedPort(uri.getPort()));
}
})
.addNodeAddress("redis://127.0.0.1:" + redisClusterContainer.getFirstMappedPort());
RedissonClient redisson = Redisson.create(config);
Map<String, Object> testMap = new HashMap<>();
testMap.put("a", "b");
testMap.put("c", "d");
testMap.put("e", "f");
testMap.put("g", "h");
testMap.put("i", "j");
testMap.put("k", "l");
RFunction f = redisson.getFunction();
f.flush();
f.load("lib", "redis.register_function('myfun', function(keys, args) return args[1] end)");
// waiting for the function replication to all nodes
Thread.sleep(5000);
RBatch batch = redisson.createBatch();
RFunctionAsync function = batch.getFunction();
for (Map.Entry<String, Object> property : testMap.entrySet()) {
List<Object> key = Collections.singletonList(property.getKey());
function.callAsync(
FunctionMode.READ,
"myfun",
FunctionResult.VALUE,
key,
property.getValue());
}
List<String> results = (List<String>) batch.execute().getResponses();
assertThat(results).containsExactly("b", "d", "f", "h", "j", "l");
redisson.shutdown();
redisClusterContainer.stop();
public void testCluster() {
testInCluster(r -> {
Map<String, Object> testMap = new HashMap<>();
testMap.put("a", "b");
testMap.put("c", "d");
testMap.put("e", "f");
testMap.put("g", "h");
testMap.put("i", "j");
testMap.put("k", "l");
RFunction f = redisson.getFunction();
f.flush();
f.load("lib", "redis.register_function('myfun', function(keys, args) return args[1] end)");
// waiting for the function replication to all nodes
try {
Thread.sleep(5000);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
RBatch batch = redisson.createBatch();
RFunctionAsync function = batch.getFunction();
for (Map.Entry<String, Object> property : testMap.entrySet()) {
List<Object> key = Collections.singletonList(property.getKey());
function.callAsync(
FunctionMode.READ,
"myfun",
FunctionResult.VALUE,
key,
property.getValue());
}
List<String> results = (List<String>) batch.execute().getResponses();
assertThat(results).containsExactly("b", "d", "f", "h", "j", "l");
});
}
@Test

@ -1,19 +1,16 @@
package org.redisson;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.redisson.api.*;
import org.redisson.api.geo.GeoSearchArgs;
import java.io.IOException;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
public class RedissonGeoTest extends BaseTest {
public class RedissonGeoTest extends RedisDockerTest {
@Test
public void testAdd() {
@ -37,8 +34,6 @@ public class RedissonGeoTest extends BaseTest {
@Test
public void testTryAdd() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RGeo<String> geo = redisson.getGeo("test");
assertThat(geo.add(2.51, 3.12, "city1")).isEqualTo(1);
assertThat(geo.tryAdd(2.5, 3.1, "city1")).isFalse();
@ -143,8 +138,6 @@ public class RedissonGeoTest extends BaseTest {
@Test
public void testBox() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RGeo<String> geo = redisson.getGeo("test");
geo.add(new GeoEntry(13.361389, 38.115556, "Palermo"), new GeoEntry(15.087269, 37.502669, "Catania"));
@ -156,8 +149,6 @@ public class RedissonGeoTest extends BaseTest {
@Test
public void testBoxWithDistance() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RGeo<String> geo = redisson.getGeo("test");
geo.add(new GeoEntry(13.361389, 38.115556, "Palermo"), new GeoEntry(15.087269, 37.502669, "Catania"));
@ -172,8 +163,6 @@ public class RedissonGeoTest extends BaseTest {
@Test
public void testBoxWithPosition() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RGeo<String> geo = redisson.getGeo("test");
geo.add(new GeoEntry(13.361389, 38.115556, "Palermo"), new GeoEntry(15.087269, 37.502669, "Catania"));
@ -188,8 +177,6 @@ public class RedissonGeoTest extends BaseTest {
@Test
public void testBoxStoreSearch() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RGeo<String> geoSource = redisson.getGeo("test");
RGeo<String> geoDest = redisson.getGeo("test-store");
geoSource.add(new GeoEntry(13.361389, 38.115556, "Palermo"), new GeoEntry(15.087269, 37.502669, "Catania"));
@ -202,8 +189,6 @@ public class RedissonGeoTest extends BaseTest {
@Test
public void testBoxStoreSorted() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RGeo<String> geoSource = redisson.getGeo("test");
RGeo<String> geoDest = redisson.getGeo("test-store");
geoSource.add(new GeoEntry(13.361389, 38.115556, "Palermo"), new GeoEntry(15.087269, 37.502669, "Catania"));

@ -13,7 +13,7 @@ import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
public class RedissonJsonBucketTest extends BaseTest {
public class RedissonJsonBucketTest extends DockerRedisStackTest {
public static class NestedType {

@ -14,6 +14,7 @@ public class RedissonRuntimeEnvironment {
public static final String OS;
public static final boolean isWindows;
private static final String MAC_PATH = "/usr/local/opt/redis/bin/redis-server";
// private static final String WINDOW_PATH = "C:\\redis\\redis-server2.cmd";
private static final String WINDOW_PATH = "C:\\redis\\redis-server.exe";
static {

@ -1,7 +1,6 @@
package org.redisson;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.Test;
import org.redisson.api.*;
import org.redisson.api.listener.ScoredSortedSetAddListener;
@ -25,7 +24,7 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.entry;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class RedissonScoredSortedSetTest extends BaseTest {
public class RedissonScoredSortedSetTest extends RedisDockerTest {
@Test
public void testEntries() {
@ -34,6 +33,10 @@ public class RedissonScoredSortedSetTest extends BaseTest {
set.add(1.2, "v2");
set.add(1.3, "v3");
RScoredSortedSet<String> set2 = redisson.getScoredSortedSet("test2");
ScoredEntry<String> s3 = set2.firstEntry();
assertThat(s3).isNull();
ScoredEntry<String> s = set.firstEntry();
assertThat(s).isEqualTo(new ScoredEntry<>(1.1, "v1"));
ScoredEntry<String> s2 = set.lastEntry();
@ -123,8 +126,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testRandom() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RScoredSortedSet<Integer> set = redisson.getScoredSortedSet("test");
set.add(1, 10);
set.add(2, 20);
@ -139,8 +140,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testTakeFirst() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("5.0.0") > 0);
final RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
Executors.newSingleThreadScheduledExecutor().schedule(() -> {
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
@ -156,8 +155,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollFirstFromAny() throws InterruptedException {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("5.0.0") > 0);
final RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
Executors.newSingleThreadScheduledExecutor().schedule(() -> {
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
@ -176,8 +173,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollFirstFromAnyCount() {
// Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
RScoredSortedSet<Integer> queue3 = redisson.getScoredSortedSet("queue:pollany2");
@ -202,8 +197,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollFirstEntriesFromAnyCount() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
RScoredSortedSet<Integer> queue3 = redisson.getScoredSortedSet("queue:pollany2");
@ -230,8 +223,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollLastEntriesFromAnyCount() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
RScoredSortedSet<Integer> queue3 = redisson.getScoredSortedSet("queue:pollany2");
@ -258,8 +249,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollFirstEntriesFromAnyTimeout() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
RScoredSortedSet<Integer> queue3 = redisson.getScoredSortedSet("queue:pollany2");
@ -286,8 +275,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollLastEntriesFromAnyTimeout() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
RScoredSortedSet<Integer> queue3 = redisson.getScoredSortedSet("queue:pollany2");
@ -314,8 +301,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollLastFromAnyCount() {
// Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
RScoredSortedSet<Integer> queue3 = redisson.getScoredSortedSet("queue:pollany2");
@ -340,8 +325,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollLastFromAny() throws InterruptedException {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("5.0.0") > 0);
final RScoredSortedSet<Integer> queue1 = redisson.getScoredSortedSet("queue:pollany");
Executors.newSingleThreadScheduledExecutor().schedule(() -> {
RScoredSortedSet<Integer> queue2 = redisson.getScoredSortedSet("queue:pollany1");
@ -820,8 +803,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollLastTimeout() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("5.0.0") > 0);
RScoredSortedSet<String> set = redisson.getScoredSortedSet("simple");
assertThat(set.pollLast(1, TimeUnit.SECONDS)).isNull();
@ -835,8 +816,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollFirstTimeout() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("5.0.0") > 0);
RScoredSortedSet<String> set = redisson.getScoredSortedSet("simple");
assertThat(set.pollFirst(1, TimeUnit.SECONDS)).isNull();
@ -850,8 +829,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollFirstTimeoutCount() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<String> set = redisson.getScoredSortedSet("simple");
assertThat(set.pollFirst(1, TimeUnit.SECONDS)).isNull();
@ -871,8 +848,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testPollLastTimeoutCount() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("7.0.0") > 0);
RScoredSortedSet<String> set = redisson.getScoredSortedSet("simple");
assertThat(set.pollFirst(1, TimeUnit.SECONDS)).isNull();
@ -1726,8 +1701,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testReadIntersection() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RScoredSortedSet<String> set1 = redisson.getScoredSortedSet("simple1");
set1.add(1, "one");
set1.add(2, "two");
@ -1833,8 +1806,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testRangeTo() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RScoredSortedSet<Integer> set1 = redisson.getScoredSortedSet("simple1");
for (int i = 0; i < 10; i++) {
set1.add(i, i);
@ -1848,8 +1819,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testRevRange() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RScoredSortedSet<Integer> set1 = redisson.getScoredSortedSet("simple1");
for (int i = 0; i < 10; i++) {
set1.add(i, i);
@ -1863,8 +1832,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testRangeToScored() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RScoredSortedSet<Integer> set1 = redisson.getScoredSortedSet("simple1");
for (int i = 0; i < 10; i++) {
set1.add(i, i);
@ -1878,8 +1845,6 @@ public class RedissonScoredSortedSetTest extends BaseTest {
@Test
public void testReadUnion() {
Assumptions.assumeTrue(RedisRunner.getDefaultRedisServerInstance().getRedisVersion().compareTo("6.2.0") > 0);
RScoredSortedSet<String> set1 = redisson.getScoredSortedSet("simple1");
set1.add(1, "one");
set1.add(2, "two");

@ -22,7 +22,7 @@ import java.util.stream.Collectors;
import static org.assertj.core.api.Assertions.assertThat;
public class RedissonSearchTest extends BaseTest {
public class RedissonSearchTest extends DockerRedisStackTest {
public static class SimpleObject {
@ -52,10 +52,10 @@ public class RedissonSearchTest extends BaseTest {
@Test
public void testMapAggregateWithCursor() {
RMap<String, SimpleObject> m = redisson.getMap("doc:1", new CompositeCodec(StringCodec.INSTANCE, redisson.getConfig().getCodec()));
RMap<String, Object> m = redisson.getMap("doc:1", new CompositeCodec(StringCodec.INSTANCE, redisson.getConfig().getCodec()));
m.put("t1", new SimpleObject("name1"));
m.put("t2", new SimpleObject("name2"));
RMap<String, SimpleObject> m2 = redisson.getMap("doc:2", new CompositeCodec(StringCodec.INSTANCE, redisson.getConfig().getCodec()));
RMap<String, Object> m2 = redisson.getMap("doc:2", new CompositeCodec(StringCodec.INSTANCE, redisson.getConfig().getCodec()));
m2.put("t1", new SimpleObject("name3"));
m2.put("t2", new SimpleObject("name4"));
@ -78,13 +78,14 @@ public class RedissonSearchTest extends BaseTest {
assertThat(r3.getTotal()).isEqualTo(1);
assertThat(r3.getCursorId()).isPositive();
assertThat(new HashSet<>(r3.getAttributes())).isEqualTo(new HashSet<>(Arrays.asList(m2.readAllMap())));
assertThat(r3.getAttributes()).hasSize(1).isSubsetOf(m.readAllMap(), m2.readAllMap());
AggregationResult r2 = s.readCursor("idx", r3.getCursorId());
assertThat(r2.getTotal()).isEqualTo(1);
assertThat(r2.getCursorId()).isPositive();
assertThat(new HashSet<>(r2.getAttributes())).isEqualTo(new HashSet<>(Arrays.asList(m.readAllMap())));
assertThat(r3.getAttributes()).isNotEqualTo(r2.getAttributes());
assertThat(r2.getAttributes()).hasSize(1).isSubsetOf(m.readAllMap(), m2.readAllMap());
}
@ -128,7 +129,7 @@ public class RedissonSearchTest extends BaseTest {
AggregationResult r = s.aggregate("idx", "*", AggregationOptions.defaults()
.load("t1", "t2"));
assertThat(r.getTotal()).isEqualTo(2);
assertThat(r.getTotal()).isEqualTo(1);
assertThat(r.getCursorId()).isEqualTo(-1);
assertThat(new HashSet<>(r.getAttributes())).isEqualTo(new HashSet<>(Arrays.asList(m2.readAllMap(), m.readAllMap())));
}

@ -13,10 +13,9 @@ import java.util.*;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.as;
import static org.assertj.core.api.Assertions.assertThat;
public class RedissonSetCacheTest extends BaseTest {
public class RedissonSetCacheTest extends RedisDockerTest {
public static class SimpleBean implements Serializable {
@ -175,9 +174,8 @@ public class RedissonSetCacheTest extends BaseTest {
assertThat(set).contains("123");
Thread.sleep(500);
assertThat(set.size()).isEqualTo(1);
assertThat(set).doesNotContain("123");
assertThat(set.contains("123")).isFalse();
assertThat(set.add("123", 1, TimeUnit.SECONDS)).isTrue();
set.destroy();
@ -214,10 +212,11 @@ public class RedissonSetCacheTest extends BaseTest {
public void testAddExpireThenAdd() throws InterruptedException, ExecutionException {
RSetCache<String> set = redisson.getSetCache("simple31");
assertThat(set.add("123", 500, TimeUnit.MILLISECONDS)).isTrue();
assertThat(set.size()).isEqualTo(1);
Thread.sleep(500);
assertThat(set.size()).isEqualTo(1);
assertThat(set.size()).isZero();
assertThat(set.contains("123")).isFalse();
assertThat(set.add("123")).isTrue();

@ -1,7 +1,20 @@
package org.redisson.executor;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
import mockit.Invocation;
import mockit.Mock;
import mockit.MockUp;
import org.junit.jupiter.api.*;
import org.redisson.BaseTest;
import org.redisson.RedisRunner;
import org.redisson.Redisson;
import org.redisson.RedissonNode;
import org.redisson.api.*;
import org.redisson.api.annotation.RInject;
import org.redisson.api.executor.TaskFinishedListener;
import org.redisson.api.executor.TaskStartedListener;
import org.redisson.config.Config;
import org.redisson.config.RedissonNodeConfig;
import org.redisson.connection.balancer.RandomLoadBalancer;
import java.io.IOException;
import java.io.Serializable;
@ -12,22 +25,8 @@ import java.util.List;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.redisson.*;
import org.redisson.api.*;
import org.redisson.api.annotation.RInject;
import org.redisson.api.executor.TaskFinishedListener;
import org.redisson.api.executor.TaskStartedListener;
import org.redisson.config.Config;
import org.redisson.config.RedissonNodeConfig;
import org.redisson.connection.balancer.RandomLoadBalancer;
import mockit.Invocation;
import mockit.Mock;
import mockit.MockUp;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
public class RedissonExecutorServiceTest extends BaseTest {
@ -455,6 +454,7 @@ public class RedissonExecutorServiceTest extends BaseTest {
}
@Test
@Timeout(1)
public void testRejectExecute() {
Assertions.assertThrows(RejectedExecutionException.class, () -> {
RExecutorService e = redisson.getExecutorService("test");

@ -52,7 +52,7 @@ public class RedissonBlockingDequeRxTest extends BaseRxTest {
RBlockingDequeRx<String> blockingDeque = redisson.getBlockingDeque("blocking_deque");
long start = System.currentTimeMillis();
String redisTask = sync(blockingDeque.pollLastAndOfferFirstTo("deque", 1, TimeUnit.SECONDS));
assertThat(System.currentTimeMillis() - start).isBetween(950L, 1500L);
assertThat(System.currentTimeMillis() - start).isBetween(950L, 1600L);
assertThat(redisTask).isNull();
}

@ -1,32 +0,0 @@
package org.redisson.spring.session;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.spring.session.config.EnableRedissonHttpSession;
import org.redisson.spring.session.config.EnableRedissonWebSession;
import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.web.reactive.DispatcherHandler;
import org.springframework.web.server.WebHandler;
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
@EnableRedissonHttpSession
//@EnableRedissonWebSession
public class HttpConfig {
@Bean
public RedissonClient redisson() {
return Redisson.create();
}
@Bean
public SessionEventsListener listener() {
return new SessionEventsListener();
}
@Bean(WebHttpHandlerBuilder.WEB_HANDLER_BEAN_NAME)
public WebHandler dispatcherHandler(ApplicationContext context) {
return new DispatcherHandler(context);
}
}

@ -1,21 +0,0 @@
package org.redisson.spring.session;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.spring.session.config.EnableRedissonHttpSession;
import org.springframework.context.annotation.Bean;
@EnableRedissonHttpSession(maxInactiveIntervalInSeconds = 5)
public class HttpConfigTimeout {
@Bean
public RedissonClient redisson() {
return Redisson.create();
}
@Bean
public SessionEventsListener listener() {
return new SessionEventsListener();
}
}

@ -1,16 +0,0 @@
package org.redisson.spring.session;
import org.springframework.session.web.context.AbstractHttpSessionApplicationInitializer;
public class HttpInitializer extends AbstractHttpSessionApplicationInitializer {
public static Class<?> CONFIG_CLASS = HttpConfig.class;
public HttpInitializer() {
super(CONFIG_CLASS);
}
@Override
public void onStartup(jakarta.servlet.ServletContext servletContext) {
}
}

@ -1,244 +0,0 @@
package org.redisson.spring.session;
import java.io.IOException;
import org.apache.http.client.ClientProtocolException;
import org.apache.http.client.fluent.Executor;
import org.apache.http.client.fluent.Request;
import org.apache.http.cookie.Cookie;
import org.apache.http.impl.client.BasicCookieStore;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.redisson.RedisRunner;
import org.redisson.RedisRunner.KEYSPACE_EVENTS_OPTIONS;
import org.redisson.RedissonRuntimeEnvironment;
import org.springframework.web.context.WebApplicationContext;
import org.springframework.web.context.support.WebApplicationContextUtils;
@Deprecated
public class RedissonSessionManagerTest {
private static RedisRunner.RedisProcess defaultRedisInstance;
@AfterAll
public static void afterClass() throws IOException, InterruptedException {
defaultRedisInstance.stop();
}
@BeforeAll
public static void beforeClass() throws IOException, InterruptedException {
if (!RedissonRuntimeEnvironment.isTravis) {
defaultRedisInstance = new RedisRunner()
.nosave()
.port(6379)
.randomDir()
.notifyKeyspaceEvents(KEYSPACE_EVENTS_OPTIONS.E,
KEYSPACE_EVENTS_OPTIONS.x,
KEYSPACE_EVENTS_OPTIONS.g)
.run();
}
}
@Test
public void testSwitchServer() throws Exception {
// start the server at http://localhost:8080/myapp
TomcatServer server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
Executor executor = Executor.newInstance();
BasicCookieStore cookieStore = new BasicCookieStore();
executor.use(cookieStore);
write(executor, "test", "1234");
Cookie cookie = cookieStore.getCookies().get(0);
Executor.closeIdleConnections();
server.stop();
server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
executor = Executor.newInstance();
cookieStore = new BasicCookieStore();
cookieStore.addCookie(cookie);
executor.use(cookieStore);
read(executor, "test", "1234");
remove(executor, "test", "null");
Executor.closeIdleConnections();
server.stop();
}
@Test
public void testWriteReadRemove() throws Exception {
// start the server at http://localhost:8080/myapp
TomcatServer server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
Executor executor = Executor.newInstance();
write(executor, "test", "1234");
read(executor, "test", "1234");
remove(executor, "test", "null");
Executor.closeIdleConnections();
server.stop();
}
@Test
public void testRecreate() throws Exception {
// start the server at http://localhost:8080/myapp
TomcatServer server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
Executor executor = Executor.newInstance();
write(executor, "test", "1");
recreate(executor, "test", "2");
read(executor, "test", "2");
Executor.closeIdleConnections();
server.stop();
}
@Test
public void testUpdate() throws Exception {
// start the server at http://localhost:8080/myapp
TomcatServer server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
Executor executor = Executor.newInstance();
write(executor, "test", "1");
read(executor, "test", "1");
write(executor, "test", "2");
read(executor, "test", "2");
Executor.closeIdleConnections();
server.stop();
}
@Test
public void testExpire() throws Exception {
HttpInitializer.CONFIG_CLASS = HttpConfigTimeout.class;
// start the server at http://localhost:8080/myapp
TomcatServer server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
WebApplicationContext wa = WebApplicationContextUtils.getRequiredWebApplicationContext(server.getServletContext());
SessionEventsListener listener = wa.getBean(SessionEventsListener.class);
Executor executor = Executor.newInstance();
BasicCookieStore cookieStore = new BasicCookieStore();
executor.use(cookieStore);
write(executor, "test", "1234");
Cookie cookie = cookieStore.getCookies().get(0);
Thread.sleep(50);
Assertions.assertEquals(1, listener.getSessionCreatedEvents());
Assertions.assertEquals(0, listener.getSessionExpiredEvents());
Executor.closeIdleConnections();
Thread.sleep(6000);
Assertions.assertEquals(1, listener.getSessionCreatedEvents());
Assertions.assertEquals(1, listener.getSessionExpiredEvents());
executor = Executor.newInstance();
cookieStore = new BasicCookieStore();
cookieStore.addCookie(cookie);
executor.use(cookieStore);
read(executor, "test", "null");
Thread.sleep(50);
Assertions.assertEquals(2, listener.getSessionCreatedEvents());
write(executor, "test", "1234");
Thread.sleep(3000);
read(executor, "test", "1234");
Thread.sleep(3000);
Assertions.assertEquals(1, listener.getSessionExpiredEvents());
Thread.sleep(1000);
Assertions.assertEquals(1, listener.getSessionExpiredEvents());
Thread.sleep(3000);
Assertions.assertEquals(2, listener.getSessionExpiredEvents());
Executor.closeIdleConnections();
server.stop();
}
@Test
public void testInvalidate() throws Exception {
// start the server at http://localhost:8080/myapp
TomcatServer server = new TomcatServer("myapp", 8080, "src/test/");
server.start();
WebApplicationContext wa = WebApplicationContextUtils.getRequiredWebApplicationContext(server.getServletContext());
SessionEventsListener listener = wa.getBean(SessionEventsListener.class);
Executor executor = Executor.newInstance();
BasicCookieStore cookieStore = new BasicCookieStore();
executor.use(cookieStore);
write(executor, "test", "1234");
Cookie cookie = cookieStore.getCookies().get(0);
Thread.sleep(50);
Assertions.assertEquals(1, listener.getSessionCreatedEvents());
Assertions.assertEquals(0, listener.getSessionDeletedEvents());
invalidate(executor);
Assertions.assertEquals(1, listener.getSessionCreatedEvents());
Assertions.assertEquals(1, listener.getSessionDeletedEvents());
Executor.closeIdleConnections();
executor = Executor.newInstance();
cookieStore = new BasicCookieStore();
cookieStore.addCookie(cookie);
executor.use(cookieStore);
read(executor, "test", "null");
Executor.closeIdleConnections();
server.stop();
}
private void write(Executor executor, String key, String value) throws IOException, ClientProtocolException {
String url = "http://localhost:8080/myapp/write?key=" + key + "&value=" + value;
String response = executor.execute(Request.Get(url)).returnContent().asString();
Assertions.assertEquals("OK", response);
}
private void read(Executor executor, String key, String value) throws IOException, ClientProtocolException {
String url = "http://localhost:8080/myapp/read?key=" + key;
String response = executor.execute(Request.Get(url)).returnContent().asString();
Assertions.assertEquals(value, response);
}
private void remove(Executor executor, String key, String value) throws IOException, ClientProtocolException {
String url = "http://localhost:8080/myapp/remove?key=" + key;
String response = executor.execute(Request.Get(url)).returnContent().asString();
Assertions.assertEquals(value, response);
}
private void invalidate(Executor executor) throws IOException, ClientProtocolException {
String url = "http://localhost:8080/myapp/invalidate";
String response = executor.execute(Request.Get(url)).returnContent().asString();
Assertions.assertEquals("OK", response);
}
private void recreate(Executor executor, String key, String value) throws IOException, ClientProtocolException {
String url = "http://localhost:8080/myapp/recreate?key=" + key + "&value=" + value;
String response = executor.execute(Request.Get(url)).returnContent().asString();
Assertions.assertEquals("OK", response);
}
}

@ -1,40 +0,0 @@
package org.redisson.spring.session;
import org.springframework.context.ApplicationListener;
import org.springframework.session.events.AbstractSessionEvent;
import org.springframework.session.events.SessionCreatedEvent;
import org.springframework.session.events.SessionDeletedEvent;
import org.springframework.session.events.SessionExpiredEvent;
public class SessionEventsListener implements ApplicationListener<AbstractSessionEvent> {
private int sessionCreatedEvents;
private int sessionDeletedEvents;
private int sessionExpiredEvents;
@Override
public void onApplicationEvent(AbstractSessionEvent event) {
if (event instanceof SessionCreatedEvent) {
sessionCreatedEvents++;
}
if (event instanceof SessionDeletedEvent) {
sessionDeletedEvents++;
}
if (event instanceof SessionExpiredEvent) {
sessionExpiredEvents++;
}
}
public int getSessionCreatedEvents() {
return sessionCreatedEvents;
}
public int getSessionDeletedEvents() {
return sessionDeletedEvents;
}
public int getSessionExpiredEvents() {
return sessionExpiredEvents;
}
}

@ -1,96 +0,0 @@
package org.redisson.spring.session;
import jakarta.servlet.ServletException;
import jakarta.servlet.annotation.WebServlet;
import jakarta.servlet.http.HttpServlet;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import jakarta.servlet.http.HttpSession;
import java.io.IOException;
@WebServlet(name = "/testServlet", urlPatterns = "/*")
public class TestServlet extends HttpServlet {
private static final long serialVersionUID = 1243830648280853203L;
@Override
protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
HttpSession session = req.getSession();
if (req.getPathInfo().equals("/write")) {
String[] params = req.getQueryString().split("&");
String key = null;
String value = null;
for (String param : params) {
String[] paramLine = param.split("=");
String keyParam = paramLine[0];
String valueParam = paramLine[1];
if ("key".equals(keyParam)) {
key = valueParam;
}
if ("value".equals(keyParam)) {
value = valueParam;
}
}
session.setAttribute(key, value);
resp.getWriter().print("OK");
} else if (req.getPathInfo().equals("/read")) {
String[] params = req.getQueryString().split("&");
String key = null;
for (String param : params) {
String[] line = param.split("=");
String keyParam = line[0];
if ("key".equals(keyParam)) {
key = line[1];
}
}
Object attr = session.getAttribute(key);
resp.getWriter().print(attr);
} else if (req.getPathInfo().equals("/remove")) {
String[] params = req.getQueryString().split("&");
String key = null;
for (String param : params) {
String[] line = param.split("=");
String keyParam = line[0];
if ("key".equals(keyParam)) {
key = line[1];
}
}
session.removeAttribute(key);
resp.getWriter().print(String.valueOf(session.getAttribute(key)));
} else if (req.getPathInfo().equals("/invalidate")) {
session.invalidate();
resp.getWriter().print("OK");
} else if (req.getPathInfo().equals("/recreate")) {
session.invalidate();
session = req.getSession();
String[] params = req.getQueryString().split("&");
String key = null;
String value = null;
for (String param : params) {
String[] paramLine = param.split("=");
String keyParam = paramLine[0];
String valueParam = paramLine[1];
if ("key".equals(keyParam)) {
key = valueParam;
}
if ("value".equals(keyParam)) {
value = valueParam;
}
}
session.setAttribute(key, value);
resp.getWriter().print("OK");
}
}
}

@ -1,64 +0,0 @@
package org.redisson.spring.session;
import java.io.File;
import java.net.MalformedURLException;
import jakarta.servlet.ServletContext;
import jakarta.servlet.ServletException;
import org.apache.catalina.LifecycleException;
import org.apache.catalina.core.StandardContext;
import org.apache.catalina.startup.Tomcat;
import org.apache.catalina.webresources.DirResourceSet;
import org.apache.catalina.webresources.StandardRoot;
public class TomcatServer {
private Tomcat tomcat = new Tomcat();
private StandardContext ctx;
public TomcatServer(String contextPath, int port, String appBase) throws MalformedURLException, ServletException {
if(contextPath == null || appBase == null || appBase.length() == 0) {
throw new IllegalArgumentException("Context path or appbase should not be null");
}
if(!contextPath.startsWith("/")) {
contextPath = "/" + contextPath;
}
tomcat.setBaseDir("."); // location where temp dir is created
tomcat.setPort(port);
tomcat.getHost().setAppBase(".");
ctx = (StandardContext) tomcat.addWebapp(contextPath, appBase);
ctx.setDelegate(true);
File additionWebInfClasses = new File("target/test-classes");
StandardRoot resources = new StandardRoot();
DirResourceSet webResourceSet = new DirResourceSet();
webResourceSet.setBase(additionWebInfClasses.toString());
webResourceSet.setWebAppMount("/WEB-INF/classes");
resources.addPostResources(webResourceSet);
ctx.setResources(resources);
}
/**
* Start the tomcat embedded server
*/
public void start() throws LifecycleException {
tomcat.start();
}
/**
* Stop the tomcat embedded server
*/
public void stop() throws LifecycleException {
tomcat.stop();
tomcat.destroy();
tomcat.getServer().await();
}
public ServletContext getServletContext() {
return ctx.getServletContext();
}
}

@ -1,30 +0,0 @@
package org.redisson.spring.session;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.spring.session.config.EnableRedissonWebSession;
import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.web.reactive.DispatcherHandler;
import org.springframework.web.server.WebHandler;
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
@EnableRedissonWebSession
public class WebConfig {
@Bean
public RedissonClient redisson() {
return Redisson.create();
}
@Bean
public SessionEventsListener listener() {
return new SessionEventsListener();
}
@Bean(WebHttpHandlerBuilder.WEB_HANDLER_BEAN_NAME)
public WebHandler dispatcherHandler(ApplicationContext context) {
return new DispatcherHandler(context);
}
}

@ -1,14 +0,0 @@
package org.redisson.spring.session;
import org.springframework.web.server.adapter.AbstractReactiveWebInitializer;
public class WebInitializer extends AbstractReactiveWebInitializer {
public static Class<?> CONFIG_CLASS = HttpConfig.class;
@Override
protected Class<?>[] getConfigClasses() {
return new Class[] {CONFIG_CLASS};
}
}
Loading…
Cancel
Save