diff --git a/Makefile b/Makefile index ace3b52b7..aa3fd8e71 100644 --- a/Makefile +++ b/Makefile @@ -12,7 +12,7 @@ # See the License for the specific language governing permissions and # limitations under the License. -REDIS_VERSION:=3.2.0 +REDIS_VERSION:=3.2.6 SPRING_PROFILE?=ci ####### @@ -25,6 +25,8 @@ work/redis-%.conf: echo port $* >> $@ echo daemonize yes >> $@ + echo protected-mode no >> $@ + echo notify-keyspace-events Ex >> $@ echo pidfile $(shell pwd)/work/redis-$*.pid >> $@ echo logfile $(shell pwd)/work/redis-$*.log >> $@ echo save \"\" >> $@ @@ -36,6 +38,8 @@ work/redis-6379.conf: echo port 6379 >> $@ echo daemonize yes >> $@ + echo protected-mode no >> $@ + echo notify-keyspace-events Ex >> $@ echo pidfile $(shell pwd)/work/redis-6379.pid >> $@ echo logfile $(shell pwd)/work/redis-6379.log >> $@ echo save \"\" >> $@ @@ -57,6 +61,7 @@ work/sentinel-%.conf: echo port $* >> $@ echo daemonize yes >> $@ + echo protected-mode no >> $@ echo bind 0.0.0.0 >> $@ echo pidfile $(shell pwd)/work/sentinel-$*.pid >> $@ echo logfile $(shell pwd)/work/sentinel-$*.log >> $@ @@ -80,6 +85,7 @@ work/cluster-%.conf: @mkdir -p $(@D) echo port $* >> $@ + echo protected-mode no >> $@ echo cluster-enabled yes >> $@ echo cluster-config-file $(shell pwd)/work/nodes-$*.conf >> $@ echo cluster-node-timeout 5 >> $@ diff --git a/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java b/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java index 054860c29..2795a9f7b 100644 --- a/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java +++ b/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java @@ -48,12 +48,24 @@ public class RedisKeyExpiredEvent extends RedisKeyspaceEvent { /** * Creates new {@link RedisKeyExpiredEvent} - * + * * @param key * @param value */ public RedisKeyExpiredEvent(byte[] key, Object value) { - super(key); + this(null, key, value); + } + + /** + * Creates new {@link RedisKeyExpiredEvent} + * + * @pamam channel + * @param key + * @param value + * @since 1.8 + */ + public RedisKeyExpiredEvent(String channel, byte[] key, Object value) { + super(channel, key); args = ByteUtils.split(key, ':'); this.value = value; diff --git a/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java b/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java index ee136b3cc..e36162405 100644 --- a/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java +++ b/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java @@ -64,6 +64,7 @@ import org.springframework.data.redis.util.ByteUtils; import org.springframework.data.util.CloseableIterator; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; +import org.springframework.util.StringUtils; /** * Redis specific {@link KeyValueAdapter} implementation. Uses binary codec to read/write data from/to Redis. Objects @@ -367,7 +368,6 @@ public class RedisKeyValueAdapter extends AbstractKeyValueAdapter List keys = new ArrayList(ids); - if (keys.isEmpty() || keys.size() < offset) { return Collections.emptyList(); } @@ -701,28 +701,7 @@ public class RedisKeyValueAdapter extends AbstractKeyValueAdapter */ @Override public void onApplicationEvent(RedisKeyspaceEvent event) { - - LOGGER.debug("Received %s .", event); - - if (event instanceof RedisKeyExpiredEvent) { - - final RedisKeyExpiredEvent expiredEvent = (RedisKeyExpiredEvent) event; - - redisOps.execute(new RedisCallback() { - - @Override - public Void doInRedis(RedisConnection connection) throws DataAccessException { - - LOGGER.debug("Cleaning up expired key '%s' data structures in keyspace '%s'.", expiredEvent.getSource(), - expiredEvent.getKeyspace()); - - connection.sRem(toBytes(expiredEvent.getKeyspace()), expiredEvent.getId()); - new IndexWriter(connection, converter).removeKeyFromIndexes(expiredEvent.getKeyspace(), expiredEvent.getId()); - return null; - } - }); - - } + // just a customization hook } /* @@ -814,12 +793,29 @@ public class RedisKeyValueAdapter extends AbstractKeyValueAdapter if (!org.springframework.util.CollectionUtils.isEmpty(hash)) { connection.del(phantomKey); } + return hash; } }); Object value = converter.read(Object.class, new RedisData(hash)); - publishEvent(new RedisKeyExpiredEvent(key, value)); + + String channel = !ObjectUtils.isEmpty(message.getChannel()) + ? converter.getConversionService().convert(message.getChannel(), String.class) : null; + + final RedisKeyExpiredEvent event = new RedisKeyExpiredEvent(channel, key, value); + + ops.execute(new RedisCallback() { + @Override + public Void doInRedis(RedisConnection connection) throws DataAccessException { + + connection.sRem(converter.getConversionService().convert(event.getKeyspace(), byte[].class), event.getId()); + new IndexWriter(connection, converter).removeKeyFromIndexes(event.getKeyspace(), event.getId()); + return null; + } + }); + + publishEvent(event); } private boolean isKeyExpirationMessage(Message message) { diff --git a/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java b/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java index 503825903..fee229eb0 100644 --- a/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java +++ b/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java @@ -25,13 +25,28 @@ import org.springframework.context.ApplicationEvent; */ public class RedisKeyspaceEvent extends ApplicationEvent { + private final String channel; + /** * Creates new {@link RedisKeyspaceEvent}. - * + * * @param key The key that expired. Must not be {@literal null}. */ public RedisKeyspaceEvent(byte[] key) { + this(null, key); + } + + /** + * Creates new {@link RedisKeyspaceEvent}. + * + * @param channel The source channel aka subscription topic. Can be {@literal null}. + * @param key The key that expired. Must not be {@literal null}. + * @since 1.8 + */ + public RedisKeyspaceEvent(String channel, byte[] key) { + super(key); + this.channel = channel; } /* @@ -42,4 +57,13 @@ public class RedisKeyspaceEvent extends ApplicationEvent { return (byte[]) super.getSource(); } + /** + * + * @return can be {@literal null}. + * @since 1.8 + */ + public String getChannel() { + return this.channel; + } + } diff --git a/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java b/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java index 54a921e6b..ae065a466 100644 --- a/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java +++ b/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java @@ -24,6 +24,7 @@ import java.util.Date; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import org.junit.After; import org.junit.AfterClass; @@ -43,6 +44,7 @@ import org.springframework.data.redis.connection.RedisConnection; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.connection.jedis.JedisConnectionFactory; import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; +import org.springframework.data.redis.core.RedisKeyValueAdapter.EnableKeyspaceEvents; import org.springframework.data.redis.core.convert.Bucket; import org.springframework.data.redis.core.convert.KeyspaceConfiguration; import org.springframework.data.redis.core.convert.MappingConfiguration; @@ -92,6 +94,7 @@ public class RedisKeyValueAdapterTests { mappingContext.afterPropertiesSet(); adapter = new RedisKeyValueAdapter(template, mappingContext); + adapter.setEnableKeyspaceEvents(EnableKeyspaceEvents.ON_STARTUP); adapter.afterPropertiesSet(); template.execute(new RedisCallback() { @@ -239,26 +242,67 @@ public class RedisKeyValueAdapterTests { assertThat(template.opsForSet().members("persons:firstname:rand"), not(hasItem("1"))); } - @Test // DATAREDIS-425 - public void keyExpiredEventShouldRemoveHelperStructures() { + /** + * @see DATAREDIS-425 + */ + @Test + public void keyExpiredEventShouldRemoveHelperStructures() throws InterruptedException { Map map = new LinkedHashMap(); map.put("_class", Person.class.getName()); map.put("firstname", "rand"); map.put("address.country", "Andor"); + template.opsForHash().putAll("persons:1", map); + template.expire("persons:1", 1, TimeUnit.SECONDS); + template.opsForSet().add("persons", "1"); template.opsForSet().add("persons:firstname:rand", "1"); template.opsForSet().add("persons:1:idx", "persons:firstname:rand"); - adapter.onApplicationEvent(new RedisKeyExpiredEvent("persons:1".getBytes(Bucket.CHARSET))); + int iterationCount = 0; + while (template.hasKey("persons:1") && iterationCount++ < 3) { // ci might be a little slow + Thread.sleep(2000); + } + assertThat(template.hasKey("persons:1"), is(false)); assertThat(template.hasKey("persons:firstname:rand"), is(false)); assertThat(template.hasKey("persons:1:idx"), is(false)); assertThat(template.opsForSet().members("persons"), not(hasItem("1"))); } - @Test // DATAREDIS-512 + /** + * @see DATAREDIS-589 + */ + @Test + public void keyExpiredEventWithoutKeyspaceShouldBeIgnored() throws InterruptedException { + + Map map = new LinkedHashMap(); + map.put("_class", Person.class.getName()); + map.put("firstname", "rand"); + map.put("address.country", "Andor"); + + template.opsForHash().putAll("persons:1", map); + template.opsForHash().putAll("1", map); + + template.expire("1", 1, TimeUnit.SECONDS); + + template.opsForSet().add("persons", "1"); + template.opsForSet().add("persons:firstname:rand", "1"); + template.opsForSet().add("persons:1:idx", "persons:firstname:rand"); + + Thread.sleep(2000); + + assertThat(template.hasKey("persons:1"), is(true)); + assertThat(template.hasKey("persons:firstname:rand"), is(true)); + assertThat(template.hasKey("persons:1:idx"), is(true)); + assertThat(template.opsForSet().members("persons"), hasItem("1")); + } + + /** + * @see DATAREDIS-512 + */ + @Test public void putWritesIndexDataCorrectly() { Person rand = new Person();