diff --git a/pom.xml b/pom.xml index b5dc2134..c07d88bc 100644 --- a/pom.xml +++ b/pom.xml @@ -428,6 +428,16 @@ + + org.apache.maven.plugins + maven-surefire-plugin + + + false + + diff --git a/src/main/java/com/moilioncircle/redis/replicator/AbstractReplicator.java b/src/main/java/com/moilioncircle/redis/replicator/AbstractReplicator.java index ab6c8325..5e4361a8 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/AbstractReplicator.java +++ b/src/main/java/com/moilioncircle/redis/replicator/AbstractReplicator.java @@ -347,6 +347,8 @@ public void builtInCommandParserRegister() { addCommandParser(CommandName.name("XDELEX"), new XDelExParser()); // since redis 8.4 addCommandParser(CommandName.name("MSETEX"), new MSetExParser()); + // flavor-specific parsers (e.g., Valkey 9 hash field TTL) + configuration.getFlavor().extendCommandParsers(this); } @Override diff --git a/src/main/java/com/moilioncircle/redis/replicator/Constants.java b/src/main/java/com/moilioncircle/redis/replicator/Constants.java index 7f1b16bc..276ac0ad 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/Constants.java +++ b/src/main/java/com/moilioncircle/redis/replicator/Constants.java @@ -98,6 +98,7 @@ private Constants() { public static final int RDB_TYPE_STREAM_LISTPACKS_2 = 19; public static final int RDB_TYPE_SET_LISTPACK = 20; /* since redis 7.2 */ public static final int RDB_TYPE_STREAM_LISTPACKS_3 = 21; /* since redis 7.2 */ + public static final int RDB_TYPE_HASH_2 = 22; /* valkey 9, hash with field-level expiration */ public static final int RDB_TYPE_HASH_METADATA = 24; /* since redis 7.4 */ public static final int RDB_TYPE_HASH_LISTPACK_EX = 25; /* since redis 7.4 */ diff --git a/src/main/java/com/moilioncircle/redis/replicator/Flavor.java b/src/main/java/com/moilioncircle/redis/replicator/Flavor.java index 00959882..7a24c4c9 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/Flavor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/Flavor.java @@ -19,6 +19,10 @@ import java.util.HashMap; import java.util.Map; +import com.moilioncircle.redis.replicator.cmd.CommandName; +import com.moilioncircle.redis.replicator.cmd.parser.HExpireAtParser; +import com.moilioncircle.redis.replicator.cmd.parser.HExpireParser; +import com.moilioncircle.redis.replicator.cmd.parser.HPExpireParser; import com.moilioncircle.redis.replicator.rdb.DefaultRdbVisitor; import com.moilioncircle.redis.replicator.rdb.RdbVisitor; @@ -68,6 +72,14 @@ public boolean isValidRdbVersion(int version) { public RdbVisitor rdbVisitor(Replicator replicator) { return new DefaultRdbVisitor(replicator); } + + @Override + public void extendCommandParsers(Replicator replicator) { + // since valkey 9 + replicator.addCommandParser(CommandName.name("HEXPIRE"), new HExpireParser()); + replicator.addCommandParser(CommandName.name("HPEXPIRE"), new HPExpireParser()); + replicator.addCommandParser(CommandName.name("HEXPIREAT"), new HExpireAtParser()); + } }; public static Flavor toFlavor(String flavor) { diff --git a/src/main/java/com/moilioncircle/redis/replicator/FlavorSupport.java b/src/main/java/com/moilioncircle/redis/replicator/FlavorSupport.java index 5b076404..29c51556 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/FlavorSupport.java +++ b/src/main/java/com/moilioncircle/redis/replicator/FlavorSupport.java @@ -30,7 +30,10 @@ public interface FlavorSupport { boolean isValidRdbVersion(int version); RdbVisitor rdbVisitor(Replicator replicator); - + + default void extendCommandParsers(Replicator replicator) { + } + default String prepend(String suffix) { return magic().toLowerCase() + suffix; } diff --git a/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HExpireAtCommand.java b/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HExpireAtCommand.java new file mode 100644 index 00000000..177c109a --- /dev/null +++ b/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HExpireAtCommand.java @@ -0,0 +1,80 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.impl; + +import com.moilioncircle.redis.replicator.cmd.CommandSpec; + +/** + * @author otheng03 + * @since 3.12.0 + */ +@CommandSpec(command = "HEXPIREAT") +public class HExpireAtCommand extends GenericKeyCommand { + + private static final long serialVersionUID = 1L; + + private long ex; + + private byte[][] fields; + + private ExistType existType; + + private CompareType compareType; + + public HExpireAtCommand() { + } + + public HExpireAtCommand(byte[] key, byte[][] fields, long ex, ExistType existType, CompareType compareType) { + super(key); + this.fields = fields; + this.ex = ex; + this.existType = existType; + this.compareType = compareType; + } + + public long getEx() { + return ex; + } + + public void setEx(long ex) { + this.ex = ex; + } + + public byte[][] getFields() { + return fields; + } + + public void setFields(byte[][] fields) { + this.fields = fields; + } + + public ExistType getExistType() { + return existType; + } + + public void setExistType(ExistType existType) { + this.existType = existType; + } + + public CompareType getCompareType() { + return compareType; + } + + public void setCompareType(CompareType compareType) { + this.compareType = compareType; + } +} diff --git a/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HExpireCommand.java b/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HExpireCommand.java new file mode 100644 index 00000000..dfb4f359 --- /dev/null +++ b/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HExpireCommand.java @@ -0,0 +1,80 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.impl; + +import com.moilioncircle.redis.replicator.cmd.CommandSpec; + +/** + * @author otheng03 + * @since 3.12.0 + */ +@CommandSpec(command = "HEXPIRE") +public class HExpireCommand extends GenericKeyCommand { + + private static final long serialVersionUID = 1L; + + private long ex; + + private byte[][] fields; + + private ExistType existType; + + private CompareType compareType; + + public HExpireCommand() { + } + + public HExpireCommand(byte[] key, byte[][] fields, long ex, ExistType existType, CompareType compareType) { + super(key); + this.fields = fields; + this.ex = ex; + this.existType = existType; + this.compareType = compareType; + } + + public long getEx() { + return ex; + } + + public void setEx(long ex) { + this.ex = ex; + } + + public byte[][] getFields() { + return fields; + } + + public void setFields(byte[][] fields) { + this.fields = fields; + } + + public ExistType getExistType() { + return existType; + } + + public void setExistType(ExistType existType) { + this.existType = existType; + } + + public CompareType getCompareType() { + return compareType; + } + + public void setCompareType(CompareType compareType) { + this.compareType = compareType; + } +} diff --git a/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HPExpireCommand.java b/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HPExpireCommand.java new file mode 100644 index 00000000..ba3753f0 --- /dev/null +++ b/src/main/java/com/moilioncircle/redis/replicator/cmd/impl/HPExpireCommand.java @@ -0,0 +1,80 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.impl; + +import com.moilioncircle.redis.replicator.cmd.CommandSpec; + +/** + * @author otheng03 + * @since 3.12.0 + */ +@CommandSpec(command = "HPEXPIRE") +public class HPExpireCommand extends GenericKeyCommand { + + private static final long serialVersionUID = 1L; + + private long ex; + + private byte[][] fields; + + private ExistType existType; + + private CompareType compareType; + + public HPExpireCommand() { + } + + public HPExpireCommand(byte[] key, byte[][] fields, long ex, ExistType existType, CompareType compareType) { + super(key); + this.fields = fields; + this.ex = ex; + this.existType = existType; + this.compareType = compareType; + } + + public long getEx() { + return ex; + } + + public void setEx(long ex) { + this.ex = ex; + } + + public byte[][] getFields() { + return fields; + } + + public void setFields(byte[][] fields) { + this.fields = fields; + } + + public ExistType getExistType() { + return existType; + } + + public void setExistType(ExistType existType) { + this.existType = existType; + } + + public CompareType getCompareType() { + return compareType; + } + + public void setCompareType(CompareType compareType) { + this.compareType = compareType; + } +} diff --git a/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireAtParser.java b/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireAtParser.java new file mode 100644 index 00000000..4dd1139e --- /dev/null +++ b/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireAtParser.java @@ -0,0 +1,68 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.parser; + +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toBytes; +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toLong; +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toRune; +import static com.moilioncircle.redis.replicator.util.Strings.isEquals; + +import com.moilioncircle.redis.replicator.cmd.CommandParser; +import com.moilioncircle.redis.replicator.cmd.impl.CompareType; +import com.moilioncircle.redis.replicator.cmd.impl.ExistType; +import com.moilioncircle.redis.replicator.cmd.impl.HExpireAtCommand; + +/** + * @author otheng03 + * @since 3.12.0 + */ +public class HExpireAtParser implements CommandParser { + + @Override + public HExpireAtCommand parse(Object[] command) { + int idx = 1; + byte[] key = toBytes(command[idx]); + idx++; + long ex = toLong(command[idx++]); + + ExistType existType = ExistType.NONE; + CompareType compareType = CompareType.NONE; + while (idx < command.length) { + String param = toRune(command[idx]); + if (isEquals(param, "NX")) { + existType = ExistType.NX; + } else if (isEquals(param, "XX")) { + existType = ExistType.XX; + } else if (isEquals(param, "GT")) { + compareType = CompareType.GT; + } else if (isEquals(param, "LT")) { + compareType = CompareType.LT; + } else if (isEquals(param, "FIELDS")) { + break; + } + idx++; + } + + idx += 2; // skip FIELDS numFields + byte[][] fields = new byte[command.length - idx][]; + for (int i = idx, j = 0; i < command.length; i++, j++) { + fields[j] = toBytes(command[i]); + } + return new HExpireAtCommand(key, fields, ex, existType, compareType); + } + +} diff --git a/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireParser.java b/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireParser.java new file mode 100644 index 00000000..307b5cff --- /dev/null +++ b/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireParser.java @@ -0,0 +1,68 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.parser; + +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toBytes; +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toLong; +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toRune; +import static com.moilioncircle.redis.replicator.util.Strings.isEquals; + +import com.moilioncircle.redis.replicator.cmd.CommandParser; +import com.moilioncircle.redis.replicator.cmd.impl.CompareType; +import com.moilioncircle.redis.replicator.cmd.impl.ExistType; +import com.moilioncircle.redis.replicator.cmd.impl.HExpireCommand; + +/** + * @author otheng03 + * @since 3.12.0 + */ +public class HExpireParser implements CommandParser { + + @Override + public HExpireCommand parse(Object[] command) { + int idx = 1; + byte[] key = toBytes(command[idx]); + idx++; + long ex = toLong(command[idx++]); + + ExistType existType = ExistType.NONE; + CompareType compareType = CompareType.NONE; + while (idx < command.length) { + String param = toRune(command[idx]); + if (isEquals(param, "NX")) { + existType = ExistType.NX; + } else if (isEquals(param, "XX")) { + existType = ExistType.XX; + } else if (isEquals(param, "GT")) { + compareType = CompareType.GT; + } else if (isEquals(param, "LT")) { + compareType = CompareType.LT; + } else if (isEquals(param, "FIELDS")) { + break; + } + idx++; + } + + idx += 2; // skip FIELDS numFields + byte[][] fields = new byte[command.length - idx][]; + for (int i = idx, j = 0; i < command.length; i++, j++) { + fields[j] = toBytes(command[i]); + } + return new HExpireCommand(key, fields, ex, existType, compareType); + } + +} diff --git a/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HPExpireParser.java b/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HPExpireParser.java new file mode 100644 index 00000000..da96d376 --- /dev/null +++ b/src/main/java/com/moilioncircle/redis/replicator/cmd/parser/HPExpireParser.java @@ -0,0 +1,68 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.parser; + +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toBytes; +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toLong; +import static com.moilioncircle.redis.replicator.cmd.CommandParsers.toRune; +import static com.moilioncircle.redis.replicator.util.Strings.isEquals; + +import com.moilioncircle.redis.replicator.cmd.CommandParser; +import com.moilioncircle.redis.replicator.cmd.impl.CompareType; +import com.moilioncircle.redis.replicator.cmd.impl.ExistType; +import com.moilioncircle.redis.replicator.cmd.impl.HPExpireCommand; + +/** + * @author otheng03 + * @since 3.12.0 + */ +public class HPExpireParser implements CommandParser { + + @Override + public HPExpireCommand parse(Object[] command) { + int idx = 1; + byte[] key = toBytes(command[idx]); + idx++; + long ex = toLong(command[idx++]); + + ExistType existType = ExistType.NONE; + CompareType compareType = CompareType.NONE; + while (idx < command.length) { + String param = toRune(command[idx]); + if (isEquals(param, "NX")) { + existType = ExistType.NX; + } else if (isEquals(param, "XX")) { + existType = ExistType.XX; + } else if (isEquals(param, "GT")) { + compareType = CompareType.GT; + } else if (isEquals(param, "LT")) { + compareType = CompareType.LT; + } else if (isEquals(param, "FIELDS")) { + break; + } + idx++; + } + + idx += 2; // skip FIELDS numFields + byte[][] fields = new byte[command.length - idx][]; + for (int i = idx, j = 0; i < command.length; i++, j++) { + fields[j] = toBytes(command[i]); + } + return new HPExpireCommand(key, fields, ex, existType, compareType); + } + +} diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbValueVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbValueVisitor.java index 8c9ce20c..a67b95c2 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbValueVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbValueVisitor.java @@ -442,6 +442,21 @@ public T applyListQuickList2(RedisInputStream in, int version) throws IOExce return (T) list; } + @Override + public T applyHash2(RedisInputStream in, int version) throws IOException { + BaseRdbParser parser = new BaseRdbParser(in); + long len = parser.rdbLoadLen().len; + TTLByteArrayMap map = new TTLByteArrayMap(); + while (len > 0) { + byte[] field = parser.rdbLoadEncodedStringObject().first(); + byte[] value = parser.rdbLoadEncodedStringObject().first(); + long expiry = parser.rdbLoadMillisecondTime(); + map.put(field, new TTLValue(expiry == -1L ? null : expiry, value)); + len--; + } + return (T) map; + } + @Override public T applyHashMetadata(RedisInputStream in, int version) throws IOException{ BaseRdbParser parser = new BaseRdbParser(in); diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbVisitor.java index 6936fbde..bcffa69c 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/DefaultRdbVisitor.java @@ -19,6 +19,7 @@ import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FREQ; import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_IDLE; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -496,12 +497,24 @@ public Event applyListQuickList2(RedisInputStream in, int version, ContextKeyVal return context.valueOf(o18); } + @Override + public Event applyHash2(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException { + BaseRdbParser parser = new BaseRdbParser(in); + KeyValuePair> o22 = new KeyStringValueTTLHash(); + byte[] key = parser.rdbLoadEncodedStringObject().first(); + TTLByteArrayMap map = valueVisitor.applyHash2(in, version); + o22.setValueRdbType(RDB_TYPE_HASH_2); + o22.setValue(map); + o22.setKey(key); + return context.valueOf(o22); + } + @Override public Event applyHashMetadata(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException{ BaseRdbParser parser = new BaseRdbParser(in); KeyValuePair> o24 = new KeyStringValueTTLHash(); byte[] key = parser.rdbLoadEncodedStringObject().first(); - + TTLByteArrayMap map = valueVisitor.applyHashMetadata(in, version); o24.setValueRdbType(RDB_TYPE_HASH_METADATA); o24.setValue(map); @@ -646,7 +659,9 @@ protected ModuleParser lookupModuleParser(String moduleName, i return (KeyValuePair) applyStreamListPacks2(in, version, context); case RDB_TYPE_STREAM_LISTPACKS_3: return (KeyValuePair) applyStreamListPacks3(in, version, context); - case RDB_TYPE_HASH_LISTPACK_EX: + case RDB_TYPE_HASH_2: + return (KeyValuePair) applyHash2(in, version, context); + case RDB_TYPE_HASH_LISTPACK_EX: return (KeyValuePair) applyHashListPackEx(in, version, context); case RDB_TYPE_HASH_METADATA: return (KeyValuePair) applyHashMetadata(in, version, context); diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbParser.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbParser.java index ffea34d8..a36f3868 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbParser.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbParser.java @@ -29,6 +29,7 @@ import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_SELECTDB; import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_SLOT_INFO; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -296,6 +297,9 @@ public long parse() throws IOException { case RDB_TYPE_STREAM_LISTPACKS_3: event = rdbVisitor.applyStreamListPacks3(in, version, kv); break; + case RDB_TYPE_HASH_2: + event = rdbVisitor.applyHash2(in, version, kv); + break; case RDB_TYPE_HASH_METADATA: event = rdbVisitor.applyHashMetadata(in, version, kv); break; diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbValueVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbValueVisitor.java index a1c95384..ee041cbd 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbValueVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbValueVisitor.java @@ -118,10 +118,14 @@ public T applyStreamListPacks3(RedisInputStream in, int version) throws IOEx throw new UnsupportedOperationException("must implement this method."); } + public T applyHash2(RedisInputStream in, int version) throws IOException { + throw new UnsupportedOperationException("must implement this method."); + } + public T applyHashMetadata(RedisInputStream in, int version) throws IOException{ throw new UnsupportedOperationException("must implement this method."); } - + public T applyHashListPackEx(RedisInputStream in, int version) throws IOException { throw new UnsupportedOperationException("must implement this method."); } diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbVisitor.java index 0f6c7794..f9a53fe5 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/RdbVisitor.java @@ -195,10 +195,14 @@ public Event applyStreamListPacks3(RedisInputStream in, int version, ContextKeyV throw new UnsupportedOperationException("must implement this method."); } + public Event applyHash2(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException { + throw new UnsupportedOperationException("must implement this method."); + } + public Event applyHashMetadata(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException{ throw new UnsupportedOperationException("must implement this method."); } - + public Event applyHashListPackEx(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException { throw new UnsupportedOperationException("must implement this method."); } diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbValueVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbValueVisitor.java index d2f48cdf..acb5427e 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbValueVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbValueVisitor.java @@ -24,6 +24,7 @@ import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FUNCTION; import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FUNCTION2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -561,6 +562,46 @@ public T applyListQuickList2(RedisInputStream in, int version) throws IOExce } } + @Override + public T applyHash2(RedisInputStream in, int version) throws IOException { + if (this.version != -1 && this.version < 80 /* since valkey rdb version 80 */) { + // downgrade to RDB_TYPE_HASH, dropping per-field TTLs + BaseRdbParser parser = new BaseRdbParser(in); + BaseRdbEncoder encoder = new BaseRdbEncoder(); + DefaultRawByteListener listener = new DefaultRawByteListener((byte) RDB_TYPE_HASH, version); + try (ByteBufferOutputStream out = new ByteBufferOutputStream(size)) { + long len = parser.rdbLoadLen().len; + listener.handle(encoder.rdbSaveLen(len)); + while (len > 0) { + ByteArray field = parser.rdbLoadEncodedStringObject(); + encoder.rdbGenericSaveStringObject(field, out); + ByteArray value = parser.rdbLoadEncodedStringObject(); + encoder.rdbGenericSaveStringObject(value, out); + parser.rdbLoadMillisecondTime(); // ignore per-field expiry + len--; + } + listener.handle(out.toByteBuffer()); + return (T) listener.getBytes(); + } + } else { + DefaultRawByteListener listener = new DefaultRawByteListener((byte) RDB_TYPE_HASH_2, version); + replicator.addRawByteListener(listener); + try { + SkipRdbParser skipParser = new SkipRdbParser(in); + long len = skipParser.rdbLoadLen().len; + while (len > 0) { + skipParser.rdbLoadEncodedStringObject(); + skipParser.rdbLoadEncodedStringObject(); + skipParser.rdbLoadMillisecondTime(); + len--; + } + } finally { + replicator.removeRawByteListener(listener); + } + return (T) listener.getBytes(); + } + } + @Override public T applyHashMetadata(RedisInputStream in, int version) throws IOException { if (this.version != -1 && this.version < 12 /* since redis rdb version 12 */) { diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitor.java index 6deee1a4..a159198b 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitor.java @@ -17,6 +17,7 @@ package com.moilioncircle.redis.replicator.rdb.dump; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -305,6 +306,21 @@ public Event applyListQuickList2(RedisInputStream in, int version, ContextKeyVal return context.valueOf(o18); } + @Override + public Event applyHash2(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException { + BaseRdbParser parser = new BaseRdbParser(in); + KeyValuePair o22 = new DumpKeyValuePair(); + byte[] key = parser.rdbLoadEncodedStringObject().first(); + if (this.version != -1 && this.version < 80 /* since valkey rdb version 80 */) { + o22.setValueRdbType(RDB_TYPE_HASH); + } else { + o22.setValueRdbType(RDB_TYPE_HASH_2); + } + o22.setKey(key); + o22.setValue(valueVisitor.applyHash2(in, version)); + return context.valueOf(o22); + } + @Override public Event applyHashMetadata(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException{ BaseRdbParser parser = new BaseRdbParser(in); diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/DefaultDumpValueParser.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/DefaultDumpValueParser.java index 4d8b2a84..a0910a2b 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/DefaultDumpValueParser.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/DefaultDumpValueParser.java @@ -19,6 +19,7 @@ import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FUNCTION; import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FUNCTION2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -148,6 +149,8 @@ public Function parse(DumpFunction function) { return KeyValuePairs.stream(kv, valueVisitor.applyStreamListPacks2(in, 0)); case RDB_TYPE_STREAM_LISTPACKS_3: return KeyValuePairs.stream(kv, valueVisitor.applyStreamListPacks3(in, 0)); + case RDB_TYPE_HASH_2: + return KeyValuePairs.ttlHash(kv, valueVisitor.applyHash2(in, 0)); case RDB_TYPE_HASH_LISTPACK_EX: return KeyValuePairs.ttlHash(kv, valueVisitor.applyHashListPackEx(in, 0)); case RDB_TYPE_HASH_METADATA: diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/IterableDumpValueParser.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/IterableDumpValueParser.java index 985f974f..106cad37 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/IterableDumpValueParser.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/dump/parser/IterableDumpValueParser.java @@ -19,6 +19,7 @@ import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FUNCTION; import static com.moilioncircle.redis.replicator.Constants.RDB_OPCODE_FUNCTION2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -163,6 +164,8 @@ public Function parse(DumpFunction function) { return KeyValuePairs.stream(kv, valueVisitor.applyStreamListPacks2(in, 0)); case RDB_TYPE_STREAM_LISTPACKS_3: return KeyValuePairs.stream(kv, valueVisitor.applyStreamListPacks3(in, 0)); + case RDB_TYPE_HASH_2: + return KeyValuePairs.iterTTLHash(kv, valueVisitor.applyHash2(in, 0)); case RDB_TYPE_HASH_LISTPACK_EX: return KeyValuePairs.iterTTLHash(kv, valueVisitor.applyHashListPackEx(in, 0)); case RDB_TYPE_HASH_METADATA: diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/iterable/ValueIterableRdbVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/iterable/ValueIterableRdbVisitor.java index c3a73999..bdd47e63 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/iterable/ValueIterableRdbVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/iterable/ValueIterableRdbVisitor.java @@ -17,6 +17,7 @@ package com.moilioncircle.redis.replicator.rdb.iterable; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -249,12 +250,23 @@ public Event applyListQuickList2(RedisInputStream in, int version, ContextKeyVal return context.valueOf(o18); } + @Override + public Event applyHash2(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException { + BaseRdbParser parser = new BaseRdbParser(in); + KeyValuePair>> o22 = new KeyStringValueTTLMapEntryIterator(); + byte[] key = parser.rdbLoadEncodedStringObject().first(); + o22.setValueRdbType(RDB_TYPE_HASH_2); + o22.setKey(key); + o22.setValue(valueVisitor.applyHash2(in, version)); + return context.valueOf(o22); + } + @Override public Event applyHashMetadata(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException{ BaseRdbParser parser = new BaseRdbParser(in); KeyValuePair>> o24 = new KeyStringValueTTLMapEntryIterator(); byte[] key = parser.rdbLoadEncodedStringObject().first(); - + o24.setValueRdbType(RDB_TYPE_HASH_METADATA); o24.setKey(key); o24.setValue(valueVisitor.applyHashMetadata(in, version)); diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbValueVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbValueVisitor.java index a29e9445..961c35e6 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbValueVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbValueVisitor.java @@ -199,6 +199,19 @@ public T applyListQuickList2(RedisInputStream in, int version) throws IOExce return null; } + @Override + public T applyHash2(RedisInputStream in, int version) throws IOException { + SkipRdbParser skip = new SkipRdbParser(in); + long len = skip.rdbLoadLen().len; + while (len > 0) { + skip.rdbLoadEncodedStringObject(); + skip.rdbLoadEncodedStringObject(); + skip.rdbLoadMillisecondTime(); // expiry per field + len--; + } + return null; + } + @Override public T applyHashMetadata(RedisInputStream in, int version) throws IOException{ SkipRdbParser skip = new SkipRdbParser(in); diff --git a/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbVisitor.java b/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbVisitor.java index 91866437..42227dc1 100644 --- a/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbVisitor.java +++ b/src/main/java/com/moilioncircle/redis/replicator/rdb/skip/SkipRdbVisitor.java @@ -273,6 +273,14 @@ public Event applyListQuickList2(RedisInputStream in, int version, ContextKeyVal return null; } + @Override + public Event applyHash2(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException { + SkipRdbParser parser = new SkipRdbParser(in); + parser.rdbLoadEncodedStringObject(); + valueVisitor.applyHash2(in, version); + return null; + } + @Override public Event applyHashMetadata(RedisInputStream in, int version, ContextKeyValuePair context) throws IOException{ SkipRdbParser parser = new SkipRdbParser(in); diff --git a/src/test/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireParserTest.java b/src/test/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireParserTest.java new file mode 100644 index 00000000..3832370f --- /dev/null +++ b/src/test/java/com/moilioncircle/redis/replicator/cmd/parser/HExpireParserTest.java @@ -0,0 +1,108 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.cmd.parser; + +import org.junit.Test; + +import com.moilioncircle.redis.replicator.cmd.impl.CompareType; +import com.moilioncircle.redis.replicator.cmd.impl.ExistType; +import com.moilioncircle.redis.replicator.cmd.impl.HExpireAtCommand; +import com.moilioncircle.redis.replicator.cmd.impl.HExpireCommand; +import com.moilioncircle.redis.replicator.cmd.impl.HPExpireCommand; + +import junit.framework.TestCase; + +/** + * @author otheng03 + * @since 3.12.0 + */ +public class HExpireParserTest extends AbstractParserTest { + + @Test + public void testHExpire() { + HExpireParser parser = new HExpireParser(); + + HExpireCommand cmd = parser.parse(toObjectArray("hexpire mykey 100 fields 2 f1 f2".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(100L, cmd.getEx()); + assertEquals(2, cmd.getFields().length); + assertEquals("f1", cmd.getFields()[0]); + assertEquals("f2", cmd.getFields()[1]); + TestCase.assertEquals(ExistType.NONE, cmd.getExistType()); + TestCase.assertEquals(CompareType.NONE, cmd.getCompareType()); + + cmd = parser.parse(toObjectArray("hexpire mykey 100 NX fields 1 f1".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(100L, cmd.getEx()); + assertEquals(1, cmd.getFields().length); + assertEquals("f1", cmd.getFields()[0]); + TestCase.assertEquals(ExistType.NX, cmd.getExistType()); + TestCase.assertEquals(CompareType.NONE, cmd.getCompareType()); + + cmd = parser.parse(toObjectArray("hexpire mykey 100 GT fields 1 f1".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(100L, cmd.getEx()); + TestCase.assertEquals(ExistType.NONE, cmd.getExistType()); + TestCase.assertEquals(CompareType.GT, cmd.getCompareType()); + + cmd = parser.parse(toObjectArray("hexpire mykey 100 XX LT fields 1 f1".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(100L, cmd.getEx()); + TestCase.assertEquals(ExistType.XX, cmd.getExistType()); + TestCase.assertEquals(CompareType.LT, cmd.getCompareType()); + } + + @Test + public void testHPExpire() { + HPExpireParser parser = new HPExpireParser(); + + HPExpireCommand cmd = parser.parse(toObjectArray("hpexpire mykey 100000 fields 2 f1 f2".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(100000L, cmd.getEx()); + assertEquals(2, cmd.getFields().length); + assertEquals("f1", cmd.getFields()[0]); + assertEquals("f2", cmd.getFields()[1]); + TestCase.assertEquals(ExistType.NONE, cmd.getExistType()); + TestCase.assertEquals(CompareType.NONE, cmd.getCompareType()); + + cmd = parser.parse(toObjectArray("hpexpire mykey 100000 NX fields 1 f1".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(100000L, cmd.getEx()); + TestCase.assertEquals(ExistType.NX, cmd.getExistType()); + TestCase.assertEquals(CompareType.NONE, cmd.getCompareType()); + } + + @Test + public void testHExpireAt() { + HExpireAtParser parser = new HExpireAtParser(); + + HExpireAtCommand cmd = parser.parse(toObjectArray("hexpireat mykey 1614139099 fields 2 f1 f2".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(1614139099L, cmd.getEx()); + assertEquals(2, cmd.getFields().length); + assertEquals("f1", cmd.getFields()[0]); + assertEquals("f2", cmd.getFields()[1]); + TestCase.assertEquals(ExistType.NONE, cmd.getExistType()); + TestCase.assertEquals(CompareType.NONE, cmd.getCompareType()); + + cmd = parser.parse(toObjectArray("hexpireat mykey 1614139099 GT fields 1 f1".split(" "))); + assertEquals("mykey", cmd.getKey()); + assertEquals(1614139099L, cmd.getEx()); + TestCase.assertEquals(ExistType.NONE, cmd.getExistType()); + TestCase.assertEquals(CompareType.GT, cmd.getCompareType()); + } +} diff --git a/src/test/java/com/moilioncircle/redis/replicator/online/EnabledIfValkey.java b/src/test/java/com/moilioncircle/redis/replicator/online/EnabledIfValkey.java new file mode 100644 index 00000000..552c84ce --- /dev/null +++ b/src/test/java/com/moilioncircle/redis/replicator/online/EnabledIfValkey.java @@ -0,0 +1,33 @@ +/* + * Copyright 2026 + * + * 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 com.moilioncircle.redis.replicator.online; + +import static java.lang.annotation.ElementType.METHOD; +import static java.lang.annotation.ElementType.TYPE; +import static java.lang.annotation.RetentionPolicy.RUNTIME; + +import java.lang.annotation.Retention; +import java.lang.annotation.Target; + +/** + * Skips a test class or method unless {@code -Dtest.flavor=valkey} is set. + * Enforced by {@link FlavorRule}, which is wired into {@link OnlineTestBase}. + */ +@Retention(RUNTIME) +@Target({TYPE, METHOD}) +public @interface EnabledIfValkey { +} diff --git a/src/test/java/com/moilioncircle/redis/replicator/online/FlavorRule.java b/src/test/java/com/moilioncircle/redis/replicator/online/FlavorRule.java new file mode 100644 index 00000000..55a1296c --- /dev/null +++ b/src/test/java/com/moilioncircle/redis/replicator/online/FlavorRule.java @@ -0,0 +1,51 @@ +/* + * Copyright 2026 + * + * 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 com.moilioncircle.redis.replicator.online; + +import org.junit.Assume; +import org.junit.rules.TestRule; +import org.junit.runner.Description; +import org.junit.runners.model.Statement; + +/** + * Enforces flavor-gating annotations like {@link EnabledIfValkey}. + * If the test class or method is annotated and {@code -Dtest.flavor=valkey} + * is not set, the test is skipped via {@link Assume#assumeTrue}. + */ +public class FlavorRule implements TestRule { + + @Override + public Statement apply(Statement base, Description description) { + return new Statement() { + @Override + public void evaluate() throws Throwable { + if (requiresValkey(description)) { + Assume.assumeTrue( + "skipped: requires -Dtest.flavor=valkey", + "valkey".equalsIgnoreCase(System.getProperty("test.flavor"))); + } + base.evaluate(); + } + }; + } + + private static boolean requiresValkey(Description description) { + if (description.getAnnotation(EnabledIfValkey.class) != null) return true; + Class testClass = description.getTestClass(); + return testClass != null && testClass.isAnnotationPresent(EnabledIfValkey.class); + } +} diff --git a/src/test/java/com/moilioncircle/redis/replicator/online/HashFieldExpireTest.java b/src/test/java/com/moilioncircle/redis/replicator/online/HashFieldExpireTest.java new file mode 100644 index 00000000..435fc87e --- /dev/null +++ b/src/test/java/com/moilioncircle/redis/replicator/online/HashFieldExpireTest.java @@ -0,0 +1,192 @@ +/* + * Copyright 2026 otheng03 + * + * 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 com.moilioncircle.redis.replicator.online; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; + +import java.io.IOException; +import java.util.concurrent.atomic.AtomicReference; + +import org.junit.Test; + +import com.moilioncircle.redis.replicator.RedisReplicator; +import com.moilioncircle.redis.replicator.Replicator; +import com.moilioncircle.redis.replicator.client.RESP2Client; +import com.moilioncircle.redis.replicator.cmd.impl.CompareType; +import com.moilioncircle.redis.replicator.cmd.impl.ExistType; +import com.moilioncircle.redis.replicator.cmd.impl.HPExpireAtCommand; +import com.moilioncircle.redis.replicator.event.Event; +import com.moilioncircle.redis.replicator.event.EventListener; +import com.moilioncircle.redis.replicator.event.PostRdbSyncEvent; +import com.moilioncircle.redis.replicator.util.Strings; + +/** + * Online tests for hash field expiry commands: HEXPIRE, HPEXPIRE, HEXPIREAT. + * + *

Valkey always propagates hash field expiry commands to replicas as + * HPEXPIREAT (absolute millisecond timestamp), regardless of which variant + * was originally issued. These tests verify that the propagated HPEXPIREAT + * command is parsed correctly, with the original option flags (NX/XX/GT/LT) + * preserved. + * + *

Requires a running Valkey 9+ server at {@code 127.0.0.1:6379}. + * Run with {@code -Dtest.flavor=valkey} to target a Valkey server. + * + * @author otheng03 + * @since 3.12.0 + */ +@EnabledIfValkey +public class HashFieldExpireTest extends OnlineTestBase { + + /** + * HEXPIRE (relative seconds, NX) is propagated as HPEXPIREAT with ExistType.NX. + */ + @Test + public void testHExpire() throws Exception { + final AtomicReference ref = new AtomicReference<>(null); + Replicator replicator = new RedisReplicator(HOST, PORT, config()); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (event instanceof PostRdbSyncEvent) { + try (RESP2Client client = new RESP2Client(HOST, PORT, config())) { + RESP2Client.Command cmd = client.newCommand(); + cmd.invoke("DEL", "hexpire_test"); + cmd.invoke("HSET", "hexpire_test", "f1", "v1", "f2", "v2"); + cmd.invoke("HEXPIRE", "hexpire_test", "300", "NX", "FIELDS", "1", "f1"); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + if (event instanceof HPExpireAtCommand) { + HPExpireAtCommand cmd = (HPExpireAtCommand) event; + if (!"hexpire_test".equals(Strings.toString(cmd.getKey()))) return; + ref.compareAndSet(null, cmd); + try { + replicator.close(); + } catch (IOException ignored) { + } + } + } + }); + replicator.open(); + + HPExpireAtCommand cmd = ref.get(); + assertNotNull(cmd); + assertEquals("hexpire_test", Strings.toString(cmd.getKey())); + assertTrue(cmd.getEx() > System.currentTimeMillis()); + assertEquals(1, cmd.getFields().length); + assertEquals("f1", Strings.toString(cmd.getFields()[0])); + assertEquals(ExistType.NX, cmd.getExistType()); + assertEquals(CompareType.NONE, cmd.getCompareType()); + } + + /** + * HPEXPIRE (relative milliseconds, GT) is propagated as HPEXPIREAT with CompareType.GT. + * + * GT requires an existing expiry to compare against, so we first set one with NX. + * Both commands arrive as HPEXPIREAT; we filter for the GT one. + */ + @Test + public void testHPExpire() throws Exception { + final AtomicReference ref = new AtomicReference<>(null); + Replicator replicator = new RedisReplicator(HOST, PORT, config()); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (event instanceof PostRdbSyncEvent) { + try (RESP2Client client = new RESP2Client(HOST, PORT, config())) { + RESP2Client.Command cmd = client.newCommand(); + cmd.invoke("DEL", "hpexpire_test"); + cmd.invoke("HSET", "hpexpire_test", "f1", "v1"); + // NX sets an initial expiry; GT then succeeds because 300000 > 100000 + cmd.invoke("HPEXPIRE", "hpexpire_test", "100000", "NX", "FIELDS", "1", "f1"); + cmd.invoke("HPEXPIRE", "hpexpire_test", "300000", "GT", "FIELDS", "1", "f1"); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + if (event instanceof HPExpireAtCommand) { + HPExpireAtCommand cmd = (HPExpireAtCommand) event; + if (!"hpexpire_test".equals(Strings.toString(cmd.getKey()))) return; + if (cmd.getCompareType() != CompareType.GT) return; + ref.compareAndSet(null, cmd); + try { + replicator.close(); + } catch (IOException ignored) { + } + } + } + }); + replicator.open(); + + HPExpireAtCommand cmd = ref.get(); + assertNotNull(cmd); + assertEquals("hpexpire_test", Strings.toString(cmd.getKey())); + assertTrue(cmd.getEx() > System.currentTimeMillis()); + assertEquals(1, cmd.getFields().length); + assertEquals("f1", Strings.toString(cmd.getFields()[0])); + assertEquals(ExistType.NONE, cmd.getExistType()); + assertEquals(CompareType.GT, cmd.getCompareType()); + } + + /** + * HEXPIREAT (absolute seconds, LT) is propagated as HPEXPIREAT with CompareType.LT. + */ + @Test + public void testHExpireAt() throws Exception { + final AtomicReference ref = new AtomicReference<>(null); + Replicator replicator = new RedisReplicator(HOST, PORT, config()); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (event instanceof PostRdbSyncEvent) { + try (RESP2Client client = new RESP2Client(HOST, PORT, config())) { + long futureSeconds = System.currentTimeMillis() / 1000 + 3600; + RESP2Client.Command cmd = client.newCommand(); + cmd.invoke("DEL", "hexpireat_test"); + cmd.invoke("HSET", "hexpireat_test", "f1", "v1", "f2", "v2"); + cmd.invoke("HEXPIREAT", "hexpireat_test", String.valueOf(futureSeconds), + "LT", "FIELDS", "2", "f1", "f2"); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + if (event instanceof HPExpireAtCommand) { + HPExpireAtCommand cmd = (HPExpireAtCommand) event; + if (!"hexpireat_test".equals(Strings.toString(cmd.getKey()))) return; + ref.compareAndSet(null, cmd); + try { + replicator.close(); + } catch (IOException ignored) { + } + } + } + }); + replicator.open(); + + HPExpireAtCommand cmd = ref.get(); + assertNotNull(cmd); + assertEquals("hexpireat_test", Strings.toString(cmd.getKey())); + assertTrue(cmd.getEx() > System.currentTimeMillis()); + assertEquals(2, cmd.getFields().length); + assertEquals(ExistType.NONE, cmd.getExistType()); + assertEquals(CompareType.LT, cmd.getCompareType()); + } +} diff --git a/src/test/java/com/moilioncircle/redis/replicator/online/OnlineTestBase.java b/src/test/java/com/moilioncircle/redis/replicator/online/OnlineTestBase.java index 260d6482..39ab4325 100644 --- a/src/test/java/com/moilioncircle/redis/replicator/online/OnlineTestBase.java +++ b/src/test/java/com/moilioncircle/redis/replicator/online/OnlineTestBase.java @@ -16,6 +16,8 @@ package com.moilioncircle.redis.replicator.online; +import org.junit.Rule; + import com.moilioncircle.redis.replicator.Configuration; import com.moilioncircle.redis.replicator.Flavor; @@ -35,6 +37,9 @@ public abstract class OnlineTestBase { protected static final Flavor FLAVOR = "valkey".equalsIgnoreCase(System.getProperty("test.flavor")) ? Flavor.VALKEY : null; + @Rule + public final FlavorRule flavorRule = new FlavorRule(); + protected Configuration config() { Configuration config = Configuration.defaultSetting().setRetries(0); if (FLAVOR != null) config.setFlavor(FLAVOR); diff --git a/src/test/java/com/moilioncircle/redis/replicator/online/RESP2ClientTest.java b/src/test/java/com/moilioncircle/redis/replicator/online/RESP2ClientTest.java index 063e1874..55b7324a 100644 --- a/src/test/java/com/moilioncircle/redis/replicator/online/RESP2ClientTest.java +++ b/src/test/java/com/moilioncircle/redis/replicator/online/RESP2ClientTest.java @@ -22,7 +22,6 @@ import org.junit.Test; -import com.moilioncircle.redis.replicator.Configuration; import com.moilioncircle.redis.replicator.client.RESP2; import com.moilioncircle.redis.replicator.client.RESP2Client; diff --git a/src/test/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitorTest.java b/src/test/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitorTest.java index 8cc5c187..9c5cef0f 100644 --- a/src/test/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitorTest.java +++ b/src/test/java/com/moilioncircle/redis/replicator/rdb/dump/DumpRdbVisitorTest.java @@ -17,6 +17,7 @@ package com.moilioncircle.redis.replicator.rdb.dump; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH; +import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_2; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_LISTPACK_EX; import static com.moilioncircle.redis.replicator.Constants.RDB_TYPE_HASH_METADATA; @@ -42,6 +43,7 @@ import com.moilioncircle.redis.replicator.Configuration; import com.moilioncircle.redis.replicator.FileType; +import com.moilioncircle.redis.replicator.Flavor; import com.moilioncircle.redis.replicator.RedisReplicator; import com.moilioncircle.redis.replicator.Replicator; import com.moilioncircle.redis.replicator.event.Event; @@ -1106,4 +1108,130 @@ public void onEvent(Replicator replicator, Event event) { assertArrayEquals(e1.getValue().getValue(), e2.getValue()); } } + + // dump-hash2.rdb was generated against Valkey 9 (RDB v80) with: + // CONFIG SET hash-max-listpack-entries 0 # force hashtable encoding so the type byte is RDB_TYPE_HASH_2 (0x16), + // # not RDB_TYPE_HASH_LISTPACK_EX + // HSET ttlhash2 field1 value1 ... field5 value5 + // HEXPIRE ttlhash2 FIELDS 1 fieldN # one HEXPIRE per field, distinct TTLs + // SAVE + // then copy

/dump.rdb to src/test/resources/dump-hash2.rdb. + @Test + @SuppressWarnings("resource") + public void testHash2Passthrough() throws IOException { + Replicator replicator = new RedisReplicator(DumpRdbVisitorTest.class. + getClassLoader().getResourceAsStream("dump-hash2.rdb") + , FileType.RDB, Configuration.defaultSetting().setFlavor(Flavor.VALKEY)); + List> expected = new ArrayList<>(); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (event instanceof KeyStringValueTTLHash) { + KeyStringValueTTLHash kv = (KeyStringValueTTLHash) event; + String key = new String(kv.getKey()); + if (key.equals("ttlhash2") && kv.getValueRdbType() == RDB_TYPE_HASH_2) { + for (Map.Entry entry: kv.getValue().entrySet()) { + expected.add(entry); + } + } + } + } + }); + replicator.open(); + + replicator = new RedisReplicator(DumpRdbVisitorTest.class. + getClassLoader().getResourceAsStream("dump-hash2.rdb") + , FileType.RDB, Configuration.defaultSetting().setFlavor(Flavor.VALKEY)); + List> actual = new ArrayList<>(); + replicator.setRdbVisitor(new DumpRdbVisitor(replicator)); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (!(event instanceof DumpKeyValuePair)) return; + DumpKeyValuePair kv = (DumpKeyValuePair) event; + String key = new String(kv.getKey()); + if (key.equals("ttlhash2") && kv.getValueRdbType() == RDB_TYPE_HASH_2) { + DumpValueParser parser = new DefaultDumpValueParser(replicator); + parser.parse(kv, new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + KeyStringValueTTLHash kv = (KeyStringValueTTLHash) event; + for (Map.Entry entry: kv.getValue().entrySet()) { + actual.add(entry); + } + } + }); + } + } + }); + replicator.open(); + + assertEquals(expected.size(), actual.size()); + for (int i = 0; i < expected.size(); i++) { + Map.Entry e1 = expected.get(i); + Map.Entry e2 = actual.get(i); + assertArrayEquals(e1.getKey(), e2.getKey()); + assertArrayEquals(e1.getValue().getValue(), e2.getValue().getValue()); + assertEquals(e1.getValue().getExpires(), e2.getValue().getExpires()); + } + } + + @Test + @SuppressWarnings("resource") + public void testHash2DowngradeToHash() throws IOException { + Replicator replicator = new RedisReplicator(DumpRdbVisitorTest.class. + getClassLoader().getResourceAsStream("dump-hash2.rdb") + , FileType.RDB, Configuration.defaultSetting().setFlavor(Flavor.VALKEY)); + List> expected = new ArrayList<>(); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (event instanceof KeyStringValueTTLHash) { + KeyStringValueTTLHash kv = (KeyStringValueTTLHash) event; + String key = new String(kv.getKey()); + if (key.equals("ttlhash2") && kv.getValueRdbType() == RDB_TYPE_HASH_2) { + for (Map.Entry entry: kv.getValue().entrySet()) { + expected.add(entry); + } + } + } + } + }); + replicator.open(); + + replicator = new RedisReplicator(DumpRdbVisitorTest.class. + getClassLoader().getResourceAsStream("dump-hash2.rdb") + , FileType.RDB, Configuration.defaultSetting().setFlavor(Flavor.VALKEY)); + List> actual = new ArrayList<>(); + replicator.setRdbVisitor(new DumpRdbVisitor(replicator, 79)); + replicator.addEventListener(new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + if (!(event instanceof DumpKeyValuePair)) return; + DumpKeyValuePair kv = (DumpKeyValuePair) event; + String key = new String(kv.getKey()); + if (key.equals("ttlhash2") && kv.getValueRdbType() == RDB_TYPE_HASH) { + DumpValueParser parser = new DefaultDumpValueParser(replicator); + parser.parse(kv, new EventListener() { + @Override + public void onEvent(Replicator replicator, Event event) { + KeyStringValueHash kv = (KeyStringValueHash) event; + for (Map.Entry entry: kv.getValue().entrySet()) { + actual.add(entry); + } + } + }); + } + } + }); + replicator.open(); + + assertEquals(expected.size(), actual.size()); + for (int i = 0; i < expected.size(); i++) { + Map.Entry e1 = expected.get(i); + Map.Entry e2 = actual.get(i); + assertArrayEquals(e1.getKey(), e2.getKey()); + assertArrayEquals(e1.getValue().getValue(), e2.getValue()); + } + } } \ No newline at end of file diff --git a/src/test/resources/dump-hash2.rdb b/src/test/resources/dump-hash2.rdb new file mode 100644 index 00000000..27733250 Binary files /dev/null and b/src/test/resources/dump-hash2.rdb differ