From 7381fe93ca60feca9ccb883325cef91ba19ff21f Mon Sep 17 00:00:00 2001 From: Julien Ruaux Date: Fri, 18 Jul 2025 16:04:37 -0700 Subject: [PATCH] feat: Added string options to search and aggregate commands --- .../RedisModulesAsyncCommandsImpl.java | 1160 ++++++++-------- .../RedisModulesReactiveCommandsImpl.java | 1164 ++++++++-------- .../api/async/RediSearchAsyncCommands.java | 4 +- .../reactive/RediSearchReactiveCommands.java | 4 +- .../api/sync/RediSearchCommands.java | 4 +- ...dulesAdvancedClusterAsyncCommandsImpl.java | 8 +- ...esAdvancedClusterReactiveCommandsImpl.java | 1209 +++++++++-------- .../search/SearchCommandBuilder.java | 555 ++++---- .../com/redis/lettucemod/ModulesTests.java | 54 +- 9 files changed, 2116 insertions(+), 2046 deletions(-) diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesAsyncCommandsImpl.java b/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesAsyncCommandsImpl.java index bf3fe2a..c389cd2 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesAsyncCommandsImpl.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesAsyncCommandsImpl.java @@ -50,583 +50,585 @@ @SuppressWarnings("unchecked") public class RedisModulesAsyncCommandsImpl extends RedisAsyncCommandsImpl - implements RedisModulesAsyncCommands { - - private final TimeSeriesCommandBuilder timeSeriesCommandBuilder; - private final SearchCommandBuilder searchCommandBuilder; - private final BloomCommandBuilder bloomCommandBuilder; - - public RedisModulesAsyncCommandsImpl(StatefulRedisModulesConnection connection, RedisCodec codec) { - super(connection, codec); - this.timeSeriesCommandBuilder = new TimeSeriesCommandBuilder<>(codec); - this.searchCommandBuilder = new SearchCommandBuilder<>(codec); - this.bloomCommandBuilder = new BloomCommandBuilder<>(codec); - } - - @Override - public StatefulRedisModulesConnection getStatefulConnection() { - return (StatefulRedisModulesConnection) super.getStatefulConnection(); - } - - @Override - public RedisFuture tsCreate(K key, CreateOptions options) { - return dispatch(timeSeriesCommandBuilder.create(key, options)); - } - - @Override - public RedisFuture tsAlter(K key, AlterOptions options) { - return dispatch(timeSeriesCommandBuilder.alter(key, options)); - } - - @Override - public RedisFuture tsAdd(K key, Sample sample) { - return dispatch(timeSeriesCommandBuilder.add(key, sample)); - } - - @Override - public RedisFuture tsAdd(K key, Sample sample, AddOptions options) { - return dispatch(timeSeriesCommandBuilder.add(key, sample, options)); - } - - @Override - public RedisFuture tsDecrby(K key, double value) { - return dispatch(timeSeriesCommandBuilder.decrby(key, value, null)); - } - - @Override - public RedisFuture tsDecrby(K key, double value, IncrbyOptions options) { - return dispatch(timeSeriesCommandBuilder.decrby(key, value, options)); - } - - @Override - public RedisFuture tsIncrby(K key, double value) { - return dispatch(timeSeriesCommandBuilder.incrby(key, value, null)); - } - - @Override - public RedisFuture tsIncrby(K key, double value, IncrbyOptions options) { - return dispatch(timeSeriesCommandBuilder.incrby(key, value, options)); - } - - @Override - public RedisFuture> tsMadd(KeySample... samples) { - return dispatch(timeSeriesCommandBuilder.madd(samples)); - } - - @Override - public RedisFuture tsCreaterule(K sourceKey, K destKey, CreateRuleOptions options) { - return dispatch(timeSeriesCommandBuilder.createRule(sourceKey, destKey, options)); - } - - @Override - public RedisFuture tsDeleterule(K sourceKey, K destKey) { - return dispatch(timeSeriesCommandBuilder.deleteRule(sourceKey, destKey)); - } - - @Override - public RedisFuture> tsRange(K key, TimeRange range) { - return dispatch(timeSeriesCommandBuilder.range(key, range)); - } - - @Override - public RedisFuture> tsRange(K key, TimeRange range, RangeOptions options) { - return dispatch(timeSeriesCommandBuilder.range(key, range, options)); - } - - @Override - public RedisFuture> tsRevrange(K key, TimeRange range) { - return dispatch(timeSeriesCommandBuilder.revrange(key, range)); - } - - @Override - public RedisFuture> tsRevrange(K key, TimeRange range, RangeOptions options) { - return dispatch(timeSeriesCommandBuilder.revrange(key, range, options)); - } - - @Override - public RedisFuture>> tsMrange(TimeRange range) { - return dispatch(timeSeriesCommandBuilder.mrange(range)); - } - - @Override - public RedisFuture>> tsMrange(TimeRange range, MRangeOptions options) { - return dispatch(timeSeriesCommandBuilder.mrange(range, options)); - } - - @Override - public RedisFuture>> tsMrevrange(TimeRange range) { - return dispatch(timeSeriesCommandBuilder.mrevrange(range)); - } - - @Override - public RedisFuture>> tsMrevrange(TimeRange range, MRangeOptions options) { - return dispatch(timeSeriesCommandBuilder.mrevrange(range, options)); - } - - @Override - public RedisFuture tsGet(K key) { - return dispatch(timeSeriesCommandBuilder.get(key)); - } - - @Override - public RedisFuture>> tsMget(MGetOptions options) { - return dispatch(timeSeriesCommandBuilder.mget(options)); - } - - @Override - public RedisFuture>> tsMget(V... filters) { - return dispatch(timeSeriesCommandBuilder.mget(filters)); - } - - @Override - public RedisFuture>> tsMgetWithLabels(V... filters) { - return dispatch(timeSeriesCommandBuilder.mgetWithLabels(filters)); - } - - @Override - public RedisFuture> tsInfo(K key) { - return dispatch(timeSeriesCommandBuilder.info(key, false)); - } - - @Override - public RedisFuture> tsInfoDebug(K key) { - return dispatch(timeSeriesCommandBuilder.info(key, true)); - } - - @Override - public RedisFuture> tsQueryIndex(V... filters) { - return dispatch(timeSeriesCommandBuilder.queryIndex(filters)); - } - - @Override - public RedisFuture tsDel(K key, TimeRange timeRange) { - return dispatch(timeSeriesCommandBuilder.tsDel(key, timeRange)); - } - - @Override - public RedisFuture ftCreate(K index, Field... fields) { - return ftCreate(index, null, fields); - } - - @Override - public RedisFuture ftCreate(K index, com.redis.lettucemod.search.CreateOptions options, - Field... fields) { - return dispatch(searchCommandBuilder.create(index, options, fields)); - } - - @Override - public RedisFuture ftDropindex(K index) { - return dispatch(searchCommandBuilder.dropIndex(index, false)); - } - - @Override - public RedisFuture ftDropindexDeleteDocs(K index) { - return dispatch(searchCommandBuilder.dropIndex(index, true)); - } - - @Override - public RedisFuture> ftInfo(K index) { - return dispatch(searchCommandBuilder.info(index)); - } - - @Override - public RedisFuture> ftSearch(K index, V query) { - return ftSearch(index, query, null); - } - - @Override - public RedisFuture> ftSearch(K index, V query, SearchOptions options) { - return dispatch(searchCommandBuilder.search(index, query, options)); - } - - @Override - public RedisFuture> ftAggregate(K index, V query) { - return ftAggregate(index, query, (AggregateOptions) null); - } - - @Override - public RedisFuture> ftAggregate(K index, V query, AggregateOptions options) { - return dispatch(searchCommandBuilder.aggregate(index, query, options)); - } - - @Override - public RedisFuture> ftAggregate(K index, V query, CursorOptions cursor) { - return ftAggregate(index, query, cursor, null); - } - - @Override - public RedisFuture> ftAggregate(K index, V query, CursorOptions cursor, - AggregateOptions options) { - return dispatch(searchCommandBuilder.aggregate(index, query, cursor, options)); - } - - @Override - public RedisFuture> ftCursorRead(K index, long cursor) { - return dispatch(searchCommandBuilder.cursorRead(index, cursor, null)); - } - - @Override - public RedisFuture> ftCursorRead(K index, long cursor, long count) { - return dispatch(searchCommandBuilder.cursorRead(index, cursor, count)); - } - - @Override - public RedisFuture ftCursorDelete(K index, long cursor) { - return dispatch(searchCommandBuilder.cursorDelete(index, cursor)); - } - - @Override - public RedisFuture ftSugadd(K key, Suggestion suggestion) { - return dispatch(searchCommandBuilder.sugadd(key, suggestion)); - } - - @Override - public RedisFuture ftSugaddIncr(K key, Suggestion suggestion) { - return dispatch(searchCommandBuilder.sugaddIncr(key, suggestion)); - } - - @Override - public RedisFuture>> ftSugget(K key, V prefix) { - return dispatch(searchCommandBuilder.sugget(key, prefix)); - } - - @Override - public RedisFuture>> ftSugget(K key, V prefix, SuggetOptions options) { - return dispatch(searchCommandBuilder.sugget(key, prefix, options)); - } - - @Override - public RedisFuture ftSugdel(K key, V string) { - return dispatch(searchCommandBuilder.sugdel(key, string)); - } - - @Override - public RedisFuture ftSuglen(K key) { - return dispatch(searchCommandBuilder.suglen(key)); - } - - @Override - public RedisFuture ftAlter(K index, Field field) { - return dispatch(searchCommandBuilder.alter(index, field)); - } - - @Override - public RedisFuture ftAliasadd(K name, K index) { - return dispatch(searchCommandBuilder.aliasAdd(name, index)); - } - - @Override - public RedisFuture ftAliasdel(K name) { - return dispatch(searchCommandBuilder.aliasDel(name)); - } - - @Override - public RedisFuture ftAliasupdate(K name, K index) { - return dispatch(searchCommandBuilder.aliasUpdate(name, index)); - } - - @Override - public RedisFuture> ftList() { - return dispatch(searchCommandBuilder.list()); - } - - @Override - public RedisFuture> ftTagvals(K index, K field) { - return dispatch(searchCommandBuilder.tagVals(index, field)); - } - - @Override - public RedisFuture ftDictadd(K dict, V... terms) { - return dispatch(searchCommandBuilder.dictadd(dict, terms)); - } - - @Override - public RedisFuture ftDictdel(K dict, V... terms) { - return dispatch(searchCommandBuilder.dictdel(dict, terms)); - } - - @Override - public RedisFuture> ftDictdump(K dict) { - return dispatch(searchCommandBuilder.dictdump(dict)); - } - - @Override - public RedisFuture bfAdd(K key, V item) { - return dispatch(bloomCommandBuilder.bfAdd(key, item)); - } - - @Override - public RedisFuture bfCard(K key) { - return dispatch(bloomCommandBuilder.bfCard(key)); - } - - @Override - public RedisFuture bfExists(K key, V item) { - return dispatch(bloomCommandBuilder.bfExists(key, item)); - } - - @Override - public RedisFuture bfInfo(K key) { - return dispatch(bloomCommandBuilder.bfInfo(key)); - } - - @Override - public RedisFuture bfInfo(K key, BloomFilterInfoType type) { - return dispatch(bloomCommandBuilder.bfInfo(key, type)); - } - - @Override - public RedisFuture> bfInsert(K key, V... items) { - return dispatch(bloomCommandBuilder.bfInsert(key, items)); - } - - @Override - public RedisFuture> bfInsert(K key, BloomFilterInsertOptions options, V... items) { - return dispatch(bloomCommandBuilder.bfInsert(key, options, items)); - } - - @Override - public RedisFuture> bfMAdd(K key, V... items) { - return dispatch(bloomCommandBuilder.bfMAdd(key, items)); - } - - @Override - public RedisFuture> bfMExists(K key, V... items) { - return dispatch(bloomCommandBuilder.bfMExists(key, items)); - } - - @Override - public RedisFuture bfReserve(K key, double errorRate, long capacity) { - return bfReserve(key, errorRate, capacity, null); - } - - @Override - public RedisFuture bfReserve(K key, double errorRate, long capacity, BloomFilterReserveOptions options) { - return dispatch(bloomCommandBuilder.bfReserve(key, errorRate, capacity, options)); - } - - @Override - public RedisFuture cfAdd(K key, V item) { - return dispatch(bloomCommandBuilder.cfAdd(key, item)); - } - - @Override - public RedisFuture cfAddNx(K key, V item) { - return dispatch(bloomCommandBuilder.cfAddNx(key, item)); - } - - @Override - public RedisFuture cfCount(K key, V item) { - return dispatch(bloomCommandBuilder.cfCount(key, item)); - } - - @Override - public RedisFuture cfDel(K key, V item) { - return dispatch(bloomCommandBuilder.cfDel(key, item)); - } - - @Override - public RedisFuture cfExists(K key, V item) { - return dispatch(bloomCommandBuilder.cfExists(key, item)); - } - - @Override - public RedisFuture cfInfo(K key) { - return dispatch(bloomCommandBuilder.cfInfo(key)); - } - - @Override - public RedisFuture> cfInsert(K key, V... items) { - return dispatch(bloomCommandBuilder.cfInsert(key, items)); - } - - @Override - public RedisFuture> cfInsert(K key, CuckooFilterInsertOptions options, V... items) { - return dispatch(bloomCommandBuilder.cfInsert(key, items, options)); - } - - @Override - public RedisFuture> cfInsertNx(K key, V... items) { - return dispatch(bloomCommandBuilder.cfInsertNx(key, items)); - } - - @Override - public RedisFuture> cfInsertNx(K key, CuckooFilterInsertOptions options, V... items) { - return dispatch(bloomCommandBuilder.cfInsertNx(key, items, options)); - } - - @Override - public RedisFuture> cfMExists(K key, V... items) { - return dispatch(bloomCommandBuilder.cfMExists(key, items)); - } - - @Override - public RedisFuture cfReserve(K key, long capacity) { - return cfReserve(key, capacity, null); - } - - @Override - public RedisFuture cfReserve(K key, long capacity, CuckooFilterReserveOptions options) { - return dispatch(bloomCommandBuilder.cfReserve(key, capacity, options)); - } - - @Override - public RedisFuture cmsIncrBy(K key, V item, long increment) { - return dispatch(bloomCommandBuilder.cmsIncrBy(key, item, increment)); - } - - @Override - public RedisFuture> cmsIncrBy(K key, LongScoredValue... itemIncrements) { - return dispatch(bloomCommandBuilder.cmsIncrBy(key, itemIncrements)); - } - - @Override - public RedisFuture cmsInitByProb(K key, double error, double probability) { - return dispatch(bloomCommandBuilder.cmsInitByProb(key, error, probability)); - } - - @Override - public RedisFuture cmsInitByDim(K key, long width, long depth) { - return dispatch(bloomCommandBuilder.cmsInitByDim(key, width, depth)); - } - - @Override - public RedisFuture> cmsQuery(K key, V... items) { - return dispatch(bloomCommandBuilder.cmsQuery(key, items)); - } - - @Override - public RedisFuture cmsMerge(K destKey, K... keys) { - return dispatch(bloomCommandBuilder.cmsMerge(destKey, keys)); - } - - @Override - public RedisFuture cmsMerge(K destKey, LongScoredValue... sourceKeyWeights) { - return dispatch(bloomCommandBuilder.cmsMerge(destKey, sourceKeyWeights)); - } - - @Override - public RedisFuture cmsInfo(K key) { - return dispatch(bloomCommandBuilder.cmsInfo(key)); - } - - @Override - public RedisFuture>> topKAdd(K key, V... items) { - return dispatch(bloomCommandBuilder.topKAdd(key, items)); - } - - @Override - public RedisFuture>> topKIncrBy(K key, LongScoredValue... itemIncrements) { - return dispatch(bloomCommandBuilder.topKIncrBy(key, itemIncrements)); - } - - @Override - public RedisFuture topKInfo(K key) { - return dispatch(bloomCommandBuilder.topKInfo(key)); - } - - @Override - public RedisFuture> topKList(K key) { - return dispatch(bloomCommandBuilder.topKList(key)); - } - - @Override - public RedisFuture>> topKListWithScores(K key) { - return dispatch(bloomCommandBuilder.topKListWithScores(key)); - } - - @Override - public RedisFuture> topKQuery(K key, V... items) { - return dispatch(bloomCommandBuilder.topKQuery(key, items)); - } - - @Override - public RedisFuture topKReserve(K key, long k) { - return dispatch(bloomCommandBuilder.topKReserve(key, k)); - } - - @Override - public RedisFuture topKReserve(K key, long k, long width, long depth, double decay) { - return dispatch(bloomCommandBuilder.topKReserve(key, k, width, depth, decay)); - } - - @Override - public RedisFuture tDigestAdd(K key, double... values) { - return dispatch(bloomCommandBuilder.tDigestAdd(key, values)); - } - - @Override - public RedisFuture> tDigestByRank(K key, long... ranks) { - return dispatch(bloomCommandBuilder.tDigestByRank(key, ranks)); - } - - @Override - public RedisFuture> tDigestByRevRank(K key, long... revRanks) { - return dispatch(bloomCommandBuilder.tDigestByRevRank(key, revRanks)); - } - - @Override - public RedisFuture> tDigestCdf(K key, double... values) { - return dispatch(bloomCommandBuilder.tDigestCdf(key, values)); - } - - @Override - public RedisFuture tDigestCreate(K key) { - return dispatch(bloomCommandBuilder.tDigestCreate(key)); - } - - @Override - public RedisFuture tDigestCreate(K key, long compression) { - return dispatch(bloomCommandBuilder.tDigestCreate(key, compression)); - } - - @Override - public RedisFuture tDigestInfo(K key) { - return dispatch(bloomCommandBuilder.tDigestInfo(key)); - } - - @Override - public RedisFuture tDigestMax(K key) { - return dispatch(bloomCommandBuilder.tDigestMax(key)); - } - - @Override - public RedisFuture tDigestMerge(K destinationKey, K... sourceKeys) { - return dispatch(bloomCommandBuilder.tDigestMerge(destinationKey, sourceKeys)); - } - - @Override - public RedisFuture tDigestMerge(K destinationKey, TDigestMergeOptions options, K... sourceKeys) { - return dispatch(bloomCommandBuilder.tDigestMerge(destinationKey, options, sourceKeys)); - } - - @Override - public RedisFuture tDigestMin(K key) { - return dispatch(bloomCommandBuilder.tDigestMin(key)); - } - - @Override - public RedisFuture> tDigestQuantile(K key, double... quantiles) { - return dispatch(bloomCommandBuilder.tDigestQuantile(key, quantiles)); - } - - @Override - public RedisFuture> tDigestRank(K key, double... values) { - return dispatch(bloomCommandBuilder.tDigestRank(key, values)); - } - - @Override - public RedisFuture tDigestReset(K key) { - return dispatch(bloomCommandBuilder.tDigestReset(key)); - } - - @Override - public RedisFuture> tDigestRevRank(K key, double... values) { - return dispatch(bloomCommandBuilder.tDigestRevRank(key, values)); - } - - @Override - public RedisFuture tDigestTrimmedMean(K key, double lowCutQuantile, double highCutQuantile) { - return dispatch(bloomCommandBuilder.tDigestTrimmedMean(key, lowCutQuantile, highCutQuantile)); - } + implements RedisModulesAsyncCommands { + + private final TimeSeriesCommandBuilder timeSeriesCommandBuilder; + + private final SearchCommandBuilder searchCommandBuilder; + + private final BloomCommandBuilder bloomCommandBuilder; + + public RedisModulesAsyncCommandsImpl(StatefulRedisModulesConnection connection, RedisCodec codec) { + super(connection, codec); + this.timeSeriesCommandBuilder = new TimeSeriesCommandBuilder<>(codec); + this.searchCommandBuilder = new SearchCommandBuilder<>(codec); + this.bloomCommandBuilder = new BloomCommandBuilder<>(codec); + } + + @Override + public StatefulRedisModulesConnection getStatefulConnection() { + return (StatefulRedisModulesConnection) super.getStatefulConnection(); + } + + @Override + public RedisFuture tsCreate(K key, CreateOptions options) { + return dispatch(timeSeriesCommandBuilder.create(key, options)); + } + + @Override + public RedisFuture tsAlter(K key, AlterOptions options) { + return dispatch(timeSeriesCommandBuilder.alter(key, options)); + } + + @Override + public RedisFuture tsAdd(K key, Sample sample) { + return dispatch(timeSeriesCommandBuilder.add(key, sample)); + } + + @Override + public RedisFuture tsAdd(K key, Sample sample, AddOptions options) { + return dispatch(timeSeriesCommandBuilder.add(key, sample, options)); + } + + @Override + public RedisFuture tsDecrby(K key, double value) { + return dispatch(timeSeriesCommandBuilder.decrby(key, value, null)); + } + + @Override + public RedisFuture tsDecrby(K key, double value, IncrbyOptions options) { + return dispatch(timeSeriesCommandBuilder.decrby(key, value, options)); + } + + @Override + public RedisFuture tsIncrby(K key, double value) { + return dispatch(timeSeriesCommandBuilder.incrby(key, value, null)); + } + + @Override + public RedisFuture tsIncrby(K key, double value, IncrbyOptions options) { + return dispatch(timeSeriesCommandBuilder.incrby(key, value, options)); + } + + @Override + public RedisFuture> tsMadd(KeySample... samples) { + return dispatch(timeSeriesCommandBuilder.madd(samples)); + } + + @Override + public RedisFuture tsCreaterule(K sourceKey, K destKey, CreateRuleOptions options) { + return dispatch(timeSeriesCommandBuilder.createRule(sourceKey, destKey, options)); + } + + @Override + public RedisFuture tsDeleterule(K sourceKey, K destKey) { + return dispatch(timeSeriesCommandBuilder.deleteRule(sourceKey, destKey)); + } + + @Override + public RedisFuture> tsRange(K key, TimeRange range) { + return dispatch(timeSeriesCommandBuilder.range(key, range)); + } + + @Override + public RedisFuture> tsRange(K key, TimeRange range, RangeOptions options) { + return dispatch(timeSeriesCommandBuilder.range(key, range, options)); + } + + @Override + public RedisFuture> tsRevrange(K key, TimeRange range) { + return dispatch(timeSeriesCommandBuilder.revrange(key, range)); + } + + @Override + public RedisFuture> tsRevrange(K key, TimeRange range, RangeOptions options) { + return dispatch(timeSeriesCommandBuilder.revrange(key, range, options)); + } + + @Override + public RedisFuture>> tsMrange(TimeRange range) { + return dispatch(timeSeriesCommandBuilder.mrange(range)); + } + + @Override + public RedisFuture>> tsMrange(TimeRange range, MRangeOptions options) { + return dispatch(timeSeriesCommandBuilder.mrange(range, options)); + } + + @Override + public RedisFuture>> tsMrevrange(TimeRange range) { + return dispatch(timeSeriesCommandBuilder.mrevrange(range)); + } + + @Override + public RedisFuture>> tsMrevrange(TimeRange range, MRangeOptions options) { + return dispatch(timeSeriesCommandBuilder.mrevrange(range, options)); + } + + @Override + public RedisFuture tsGet(K key) { + return dispatch(timeSeriesCommandBuilder.get(key)); + } + + @Override + public RedisFuture>> tsMget(MGetOptions options) { + return dispatch(timeSeriesCommandBuilder.mget(options)); + } + + @Override + public RedisFuture>> tsMget(V... filters) { + return dispatch(timeSeriesCommandBuilder.mget(filters)); + } + + @Override + public RedisFuture>> tsMgetWithLabels(V... filters) { + return dispatch(timeSeriesCommandBuilder.mgetWithLabels(filters)); + } + + @Override + public RedisFuture> tsInfo(K key) { + return dispatch(timeSeriesCommandBuilder.info(key, false)); + } + + @Override + public RedisFuture> tsInfoDebug(K key) { + return dispatch(timeSeriesCommandBuilder.info(key, true)); + } + + @Override + public RedisFuture> tsQueryIndex(V... filters) { + return dispatch(timeSeriesCommandBuilder.queryIndex(filters)); + } + + @Override + public RedisFuture tsDel(K key, TimeRange timeRange) { + return dispatch(timeSeriesCommandBuilder.tsDel(key, timeRange)); + } + + @Override + public RedisFuture ftCreate(K index, Field... fields) { + return ftCreate(index, null, fields); + } + + @Override + public RedisFuture ftCreate(K index, com.redis.lettucemod.search.CreateOptions options, Field... fields) { + return dispatch(searchCommandBuilder.create(index, options, fields)); + } + + @Override + public RedisFuture ftDropindex(K index) { + return dispatch(searchCommandBuilder.dropIndex(index, false)); + } + + @Override + public RedisFuture ftDropindexDeleteDocs(K index) { + return dispatch(searchCommandBuilder.dropIndex(index, true)); + } + + @Override + public RedisFuture> ftInfo(K index) { + return dispatch(searchCommandBuilder.info(index)); + } + + @Override + public RedisFuture> ftSearch(K index, V query, V... options) { + return dispatch(searchCommandBuilder.search(index, query, options)); + } + + @Override + public RedisFuture> ftSearch(K index, V query, SearchOptions options) { + return dispatch(searchCommandBuilder.search(index, query, options)); + } + + @Override + public RedisFuture> ftAggregate(K index, V query, V... options) { + return dispatch(searchCommandBuilder.aggregate(index, query, options)); + } + + @Override + public RedisFuture> ftAggregate(K index, V query, AggregateOptions options) { + return dispatch(searchCommandBuilder.aggregate(index, query, options)); + } + + @Override + public RedisFuture> ftAggregate(K index, V query, CursorOptions cursor) { + return ftAggregate(index, query, cursor, null); + } + + @Override + public RedisFuture> ftAggregate(K index, V query, CursorOptions cursor, + AggregateOptions options) { + return dispatch(searchCommandBuilder.aggregate(index, query, cursor, options)); + } + + @Override + public RedisFuture> ftCursorRead(K index, long cursor) { + return dispatch(searchCommandBuilder.cursorRead(index, cursor, null)); + } + + @Override + public RedisFuture> ftCursorRead(K index, long cursor, long count) { + return dispatch(searchCommandBuilder.cursorRead(index, cursor, count)); + } + + @Override + public RedisFuture ftCursorDelete(K index, long cursor) { + return dispatch(searchCommandBuilder.cursorDelete(index, cursor)); + } + + @Override + public RedisFuture ftSugadd(K key, Suggestion suggestion) { + return dispatch(searchCommandBuilder.sugadd(key, suggestion)); + } + + @Override + public RedisFuture ftSugaddIncr(K key, Suggestion suggestion) { + return dispatch(searchCommandBuilder.sugaddIncr(key, suggestion)); + } + + @Override + public RedisFuture>> ftSugget(K key, V prefix) { + return dispatch(searchCommandBuilder.sugget(key, prefix)); + } + + @Override + public RedisFuture>> ftSugget(K key, V prefix, SuggetOptions options) { + return dispatch(searchCommandBuilder.sugget(key, prefix, options)); + } + + @Override + public RedisFuture ftSugdel(K key, V string) { + return dispatch(searchCommandBuilder.sugdel(key, string)); + } + + @Override + public RedisFuture ftSuglen(K key) { + return dispatch(searchCommandBuilder.suglen(key)); + } + + @Override + public RedisFuture ftAlter(K index, Field field) { + return dispatch(searchCommandBuilder.alter(index, field)); + } + + @Override + public RedisFuture ftAliasadd(K name, K index) { + return dispatch(searchCommandBuilder.aliasAdd(name, index)); + } + + @Override + public RedisFuture ftAliasdel(K name) { + return dispatch(searchCommandBuilder.aliasDel(name)); + } + + @Override + public RedisFuture ftAliasupdate(K name, K index) { + return dispatch(searchCommandBuilder.aliasUpdate(name, index)); + } + + @Override + public RedisFuture> ftList() { + return dispatch(searchCommandBuilder.list()); + } + + @Override + public RedisFuture> ftTagvals(K index, K field) { + return dispatch(searchCommandBuilder.tagVals(index, field)); + } + + @Override + public RedisFuture ftDictadd(K dict, V... terms) { + return dispatch(searchCommandBuilder.dictadd(dict, terms)); + } + + @Override + public RedisFuture ftDictdel(K dict, V... terms) { + return dispatch(searchCommandBuilder.dictdel(dict, terms)); + } + + @Override + public RedisFuture> ftDictdump(K dict) { + return dispatch(searchCommandBuilder.dictdump(dict)); + } + + @Override + public RedisFuture bfAdd(K key, V item) { + return dispatch(bloomCommandBuilder.bfAdd(key, item)); + } + + @Override + public RedisFuture bfCard(K key) { + return dispatch(bloomCommandBuilder.bfCard(key)); + } + + @Override + public RedisFuture bfExists(K key, V item) { + return dispatch(bloomCommandBuilder.bfExists(key, item)); + } + + @Override + public RedisFuture bfInfo(K key) { + return dispatch(bloomCommandBuilder.bfInfo(key)); + } + + @Override + public RedisFuture bfInfo(K key, BloomFilterInfoType type) { + return dispatch(bloomCommandBuilder.bfInfo(key, type)); + } + + @Override + public RedisFuture> bfInsert(K key, V... items) { + return dispatch(bloomCommandBuilder.bfInsert(key, items)); + } + + @Override + public RedisFuture> bfInsert(K key, BloomFilterInsertOptions options, V... items) { + return dispatch(bloomCommandBuilder.bfInsert(key, options, items)); + } + + @Override + public RedisFuture> bfMAdd(K key, V... items) { + return dispatch(bloomCommandBuilder.bfMAdd(key, items)); + } + + @Override + public RedisFuture> bfMExists(K key, V... items) { + return dispatch(bloomCommandBuilder.bfMExists(key, items)); + } + + @Override + public RedisFuture bfReserve(K key, double errorRate, long capacity) { + return bfReserve(key, errorRate, capacity, null); + } + + @Override + public RedisFuture bfReserve(K key, double errorRate, long capacity, BloomFilterReserveOptions options) { + return dispatch(bloomCommandBuilder.bfReserve(key, errorRate, capacity, options)); + } + + @Override + public RedisFuture cfAdd(K key, V item) { + return dispatch(bloomCommandBuilder.cfAdd(key, item)); + } + + @Override + public RedisFuture cfAddNx(K key, V item) { + return dispatch(bloomCommandBuilder.cfAddNx(key, item)); + } + + @Override + public RedisFuture cfCount(K key, V item) { + return dispatch(bloomCommandBuilder.cfCount(key, item)); + } + + @Override + public RedisFuture cfDel(K key, V item) { + return dispatch(bloomCommandBuilder.cfDel(key, item)); + } + + @Override + public RedisFuture cfExists(K key, V item) { + return dispatch(bloomCommandBuilder.cfExists(key, item)); + } + + @Override + public RedisFuture cfInfo(K key) { + return dispatch(bloomCommandBuilder.cfInfo(key)); + } + + @Override + public RedisFuture> cfInsert(K key, V... items) { + return dispatch(bloomCommandBuilder.cfInsert(key, items)); + } + + @Override + public RedisFuture> cfInsert(K key, CuckooFilterInsertOptions options, V... items) { + return dispatch(bloomCommandBuilder.cfInsert(key, items, options)); + } + + @Override + public RedisFuture> cfInsertNx(K key, V... items) { + return dispatch(bloomCommandBuilder.cfInsertNx(key, items)); + } + + @Override + public RedisFuture> cfInsertNx(K key, CuckooFilterInsertOptions options, V... items) { + return dispatch(bloomCommandBuilder.cfInsertNx(key, items, options)); + } + + @Override + public RedisFuture> cfMExists(K key, V... items) { + return dispatch(bloomCommandBuilder.cfMExists(key, items)); + } + + @Override + public RedisFuture cfReserve(K key, long capacity) { + return cfReserve(key, capacity, null); + } + + @Override + public RedisFuture cfReserve(K key, long capacity, CuckooFilterReserveOptions options) { + return dispatch(bloomCommandBuilder.cfReserve(key, capacity, options)); + } + + @Override + public RedisFuture cmsIncrBy(K key, V item, long increment) { + return dispatch(bloomCommandBuilder.cmsIncrBy(key, item, increment)); + } + + @Override + public RedisFuture> cmsIncrBy(K key, LongScoredValue... itemIncrements) { + return dispatch(bloomCommandBuilder.cmsIncrBy(key, itemIncrements)); + } + + @Override + public RedisFuture cmsInitByProb(K key, double error, double probability) { + return dispatch(bloomCommandBuilder.cmsInitByProb(key, error, probability)); + } + + @Override + public RedisFuture cmsInitByDim(K key, long width, long depth) { + return dispatch(bloomCommandBuilder.cmsInitByDim(key, width, depth)); + } + + @Override + public RedisFuture> cmsQuery(K key, V... items) { + return dispatch(bloomCommandBuilder.cmsQuery(key, items)); + } + + @Override + public RedisFuture cmsMerge(K destKey, K... keys) { + return dispatch(bloomCommandBuilder.cmsMerge(destKey, keys)); + } + + @Override + public RedisFuture cmsMerge(K destKey, LongScoredValue... sourceKeyWeights) { + return dispatch(bloomCommandBuilder.cmsMerge(destKey, sourceKeyWeights)); + } + + @Override + public RedisFuture cmsInfo(K key) { + return dispatch(bloomCommandBuilder.cmsInfo(key)); + } + + @Override + public RedisFuture>> topKAdd(K key, V... items) { + return dispatch(bloomCommandBuilder.topKAdd(key, items)); + } + + @Override + public RedisFuture>> topKIncrBy(K key, LongScoredValue... itemIncrements) { + return dispatch(bloomCommandBuilder.topKIncrBy(key, itemIncrements)); + } + + @Override + public RedisFuture topKInfo(K key) { + return dispatch(bloomCommandBuilder.topKInfo(key)); + } + + @Override + public RedisFuture> topKList(K key) { + return dispatch(bloomCommandBuilder.topKList(key)); + } + + @Override + public RedisFuture>> topKListWithScores(K key) { + return dispatch(bloomCommandBuilder.topKListWithScores(key)); + } + + @Override + public RedisFuture> topKQuery(K key, V... items) { + return dispatch(bloomCommandBuilder.topKQuery(key, items)); + } + + @Override + public RedisFuture topKReserve(K key, long k) { + return dispatch(bloomCommandBuilder.topKReserve(key, k)); + } + + @Override + public RedisFuture topKReserve(K key, long k, long width, long depth, double decay) { + return dispatch(bloomCommandBuilder.topKReserve(key, k, width, depth, decay)); + } + + @Override + public RedisFuture tDigestAdd(K key, double... values) { + return dispatch(bloomCommandBuilder.tDigestAdd(key, values)); + } + + @Override + public RedisFuture> tDigestByRank(K key, long... ranks) { + return dispatch(bloomCommandBuilder.tDigestByRank(key, ranks)); + } + + @Override + public RedisFuture> tDigestByRevRank(K key, long... revRanks) { + return dispatch(bloomCommandBuilder.tDigestByRevRank(key, revRanks)); + } + + @Override + public RedisFuture> tDigestCdf(K key, double... values) { + return dispatch(bloomCommandBuilder.tDigestCdf(key, values)); + } + + @Override + public RedisFuture tDigestCreate(K key) { + return dispatch(bloomCommandBuilder.tDigestCreate(key)); + } + + @Override + public RedisFuture tDigestCreate(K key, long compression) { + return dispatch(bloomCommandBuilder.tDigestCreate(key, compression)); + } + + @Override + public RedisFuture tDigestInfo(K key) { + return dispatch(bloomCommandBuilder.tDigestInfo(key)); + } + + @Override + public RedisFuture tDigestMax(K key) { + return dispatch(bloomCommandBuilder.tDigestMax(key)); + } + + @Override + public RedisFuture tDigestMerge(K destinationKey, K... sourceKeys) { + return dispatch(bloomCommandBuilder.tDigestMerge(destinationKey, sourceKeys)); + } + + @Override + public RedisFuture tDigestMerge(K destinationKey, TDigestMergeOptions options, K... sourceKeys) { + return dispatch(bloomCommandBuilder.tDigestMerge(destinationKey, options, sourceKeys)); + } + + @Override + public RedisFuture tDigestMin(K key) { + return dispatch(bloomCommandBuilder.tDigestMin(key)); + } + + @Override + public RedisFuture> tDigestQuantile(K key, double... quantiles) { + return dispatch(bloomCommandBuilder.tDigestQuantile(key, quantiles)); + } + + @Override + public RedisFuture> tDigestRank(K key, double... values) { + return dispatch(bloomCommandBuilder.tDigestRank(key, values)); + } + + @Override + public RedisFuture tDigestReset(K key) { + return dispatch(bloomCommandBuilder.tDigestReset(key)); + } + + @Override + public RedisFuture> tDigestRevRank(K key, double... values) { + return dispatch(bloomCommandBuilder.tDigestRevRank(key, values)); + } + + @Override + public RedisFuture tDigestTrimmedMean(K key, double lowCutQuantile, double highCutQuantile) { + return dispatch(bloomCommandBuilder.tDigestTrimmedMean(key, lowCutQuantile, highCutQuantile)); + } + } diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesReactiveCommandsImpl.java b/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesReactiveCommandsImpl.java index c00bd12..3bcbd6e 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesReactiveCommandsImpl.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/RedisModulesReactiveCommandsImpl.java @@ -49,584 +49,588 @@ @SuppressWarnings("unchecked") public class RedisModulesReactiveCommandsImpl extends RedisReactiveCommandsImpl - implements RedisModulesReactiveCommands { - - private final StatefulRedisModulesConnection connection; - private final TimeSeriesCommandBuilder timeSeriesCommandBuilder; - private final SearchCommandBuilder searchCommandBuilder; - private final BloomCommandBuilder bloomCommandBuilder; - - public RedisModulesReactiveCommandsImpl(StatefulRedisModulesConnection connection, RedisCodec codec) { - super(connection, codec); - this.connection = connection; - this.timeSeriesCommandBuilder = new TimeSeriesCommandBuilder<>(codec); - this.searchCommandBuilder = new SearchCommandBuilder<>(codec); - this.bloomCommandBuilder = new BloomCommandBuilder<>(codec); - } - - @Override - public StatefulRedisModulesConnection getStatefulConnection() { - return connection; - } - - @Override - public Mono tsCreate(K key, CreateOptions options) { - return createMono(() -> timeSeriesCommandBuilder.create(key, options)); - } - - @Override - public Mono tsAlter(K key, AlterOptions options) { - return createMono(() -> timeSeriesCommandBuilder.alter(key, options)); - } - - @Override - public Mono tsAdd(K key, Sample sample) { - return createMono(() -> timeSeriesCommandBuilder.add(key, sample)); - } - - @Override - public Mono tsAdd(K key, Sample sample, AddOptions options) { - return createMono(() -> timeSeriesCommandBuilder.add(key, sample, options)); - } - - @Override - public Mono tsIncrby(K key, double value) { - return createMono(() -> timeSeriesCommandBuilder.incrby(key, value, null)); - } - - @Override - public Mono tsIncrby(K key, double value, IncrbyOptions options) { - return createMono(() -> timeSeriesCommandBuilder.incrby(key, value, options)); - } - - @Override - public Mono tsDecrby(K key, double value) { - return createMono(() -> timeSeriesCommandBuilder.decrby(key, value, null)); - } - - @Override - public Mono tsDecrby(K key, double value, IncrbyOptions options) { - return createMono(() -> timeSeriesCommandBuilder.decrby(key, value, options)); - } - - @Override - public Flux tsMadd(KeySample... samples) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.madd(samples)); - } - - @Override - public Mono tsCreaterule(K sourceKey, K destKey, CreateRuleOptions options) { - return createMono(() -> timeSeriesCommandBuilder.createRule(sourceKey, destKey, options)); - } - - @Override - public Mono tsDeleterule(K sourceKey, K destKey) { - return createMono(() -> timeSeriesCommandBuilder.deleteRule(sourceKey, destKey)); - } - - @Override - public Flux tsRange(K key, TimeRange range) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.range(key, range)); - } - - @Override - public Flux tsRange(K key, TimeRange range, RangeOptions options) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.range(key, range, options)); - } - - @Override - public Flux tsRevrange(K key, TimeRange range) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.revrange(key, range)); - } - - @Override - public Flux tsRevrange(K key, TimeRange range, RangeOptions options) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.revrange(key, range, options)); - } - - @Override - public Flux> tsMrange(TimeRange range) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrange(range)); - } - - @Override - public Flux> tsMrange(TimeRange range, MRangeOptions options) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrange(range, options)); - } - - @Override - public Flux> tsMrevrange(TimeRange range) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrevrange(range)); - } - - @Override - public Flux> tsMrevrange(TimeRange range, MRangeOptions options) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrevrange(range, options)); - } - - @Override - public Mono tsGet(K key) { - return createMono(() -> timeSeriesCommandBuilder.get(key)); - } - - @Override - public Flux> tsMget(MGetOptions options) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mget(options)); - } - - @Override - public Flux> tsMget(V... filters) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mget(filters)); - } - - @Override - public Flux> tsMgetWithLabels(V... filters) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.mgetWithLabels(filters)); - } - - @Override - public Flux tsInfo(K key) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.info(key, false)); - } - - @Override - public Flux tsInfoDebug(K key) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.info(key, true)); - } - - @Override - public Flux tsQueryIndex(V... filters) { - return createDissolvingFlux(() -> timeSeriesCommandBuilder.queryIndex(filters)); - } - - @Override - public Mono tsDel(K key, TimeRange timeRange) { - return createMono(() -> timeSeriesCommandBuilder.tsDel(key, timeRange)); - } - - @Override - public Mono ftCreate(K index, Field... fields) { - return ftCreate(index, null, fields); - } - - @Override - public Mono ftCreate(K index, com.redis.lettucemod.search.CreateOptions options, Field... fields) { - return createMono(() -> searchCommandBuilder.create(index, options, fields)); - } - - @Override - public Mono ftDropindex(K index) { - return createMono(() -> searchCommandBuilder.dropIndex(index, false)); - } - - @Override - public Mono ftDropindexDeleteDocs(K index) { - return createMono(() -> searchCommandBuilder.dropIndex(index, true)); - } - - @Override - public Flux ftInfo(K index) { - return createDissolvingFlux(() -> searchCommandBuilder.info(index)); - } - - @Override - public Mono> ftSearch(K index, V query) { - return createMono(() -> searchCommandBuilder.search(index, query, null)); - } - - @Override - public Mono> ftSearch(K index, V query, SearchOptions options) { - return createMono(() -> searchCommandBuilder.search(index, query, options)); - } - - @Override - public Mono> ftAggregate(K index, V query) { - return createMono(() -> searchCommandBuilder.aggregate(index, query, null)); - } - - @Override - public Mono> ftAggregate(K index, V query, AggregateOptions options) { - return createMono(() -> searchCommandBuilder.aggregate(index, query, options)); - } - - @Override - public Mono> ftAggregate(K index, V query, CursorOptions cursor) { - return createMono(() -> searchCommandBuilder.aggregate(index, query, cursor, null)); - } - - @Override - public Mono> ftAggregate(K index, V query, CursorOptions cursor, - AggregateOptions options) { - return createMono(() -> searchCommandBuilder.aggregate(index, query, cursor, options)); - } - - @Override - public Mono> ftCursorRead(K index, long cursor) { - return createMono(() -> searchCommandBuilder.cursorRead(index, cursor, null)); - } - - @Override - public Mono> ftCursorRead(K index, long cursor, long count) { - return createMono(() -> searchCommandBuilder.cursorRead(index, cursor, count)); - } - - @Override - public Mono ftCursorDelete(K index, long cursor) { - return createMono(() -> searchCommandBuilder.cursorDelete(index, cursor)); - } - - @Override - public Mono ftSugadd(K key, Suggestion suggestion) { - return createMono(() -> searchCommandBuilder.sugadd(key, suggestion)); - } - - @Override - public Mono ftSugaddIncr(K key, Suggestion suggestion) { - return createMono(() -> searchCommandBuilder.sugaddIncr(key, suggestion)); - } - - @Override - public Flux> ftSugget(K key, V prefix) { - return createDissolvingFlux(() -> searchCommandBuilder.sugget(key, prefix)); - } - - @Override - public Flux> ftSugget(K key, V prefix, SuggetOptions options) { - return createDissolvingFlux(() -> searchCommandBuilder.sugget(key, prefix, options)); - } - - @Override - public Mono ftSugdel(K key, V string) { - return createMono(() -> searchCommandBuilder.sugdel(key, string)); - } - - @Override - public Mono ftSuglen(K key) { - return createMono(() -> searchCommandBuilder.suglen(key)); - } - - @Override - public Mono ftAlter(K index, Field field) { - return createMono(() -> searchCommandBuilder.alter(index, field)); - } - - @Override - public Mono ftAliasadd(K name, K index) { - return createMono(() -> searchCommandBuilder.aliasAdd(name, index)); - } - - @Override - public Mono ftAliasupdate(K name, K index) { - return createMono(() -> searchCommandBuilder.aliasUpdate(name, index)); - } - - @Override - public Mono ftAliasdel(K name) { - return createMono(() -> searchCommandBuilder.aliasDel(name)); - } - - @Override - public Flux ftList() { - return createDissolvingFlux(searchCommandBuilder::list); - } - - @Override - public Flux ftTagvals(K index, K field) { - return createDissolvingFlux(() -> searchCommandBuilder.tagVals(index, field)); - } - - @Override - public Mono ftDictadd(K dict, V... terms) { - return createMono(() -> searchCommandBuilder.dictadd(dict, terms)); - } - - @Override - public Mono ftDictdel(K dict, V... terms) { - return createMono(() -> searchCommandBuilder.dictdel(dict, terms)); - } - - @Override - public Flux ftDictdump(K dict) { - return createDissolvingFlux(() -> searchCommandBuilder.dictdump(dict)); - } - - @Override - public Flux> topKAdd(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.topKAdd(key, items)); - } - - @Override - public Flux> topKIncrBy(K key, LongScoredValue... itemIncrements) { - return createDissolvingFlux(() -> bloomCommandBuilder.topKIncrBy(key, itemIncrements)); - } - - @Override - public Mono topKInfo(K key) { - return createMono(() -> bloomCommandBuilder.topKInfo(key)); - } - - @Override - public Flux topKList(K key) { - return createDissolvingFlux(() -> bloomCommandBuilder.topKList(key)); - } - - @Override - public Flux> topKListWithScores(K key) { - return createDissolvingFlux(() -> bloomCommandBuilder.topKListWithScores(key)); - } - - @Override - public Flux topKQuery(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.topKQuery(key, items)); - } - - @Override - public Mono topKReserve(K key, long k) { - return createMono(() -> bloomCommandBuilder.topKReserve(key, k)); - } - - @Override - public Mono topKReserve(K key, long k, long width, long depth, double decay) { - return createMono(() -> bloomCommandBuilder.topKReserve(key, k, width, depth, decay)); - } - - @Override - public Mono tDigestAdd(K key, double... values) { - return createMono(() -> bloomCommandBuilder.tDigestAdd(key, values)); - } - - @Override - public Flux tDigestByRank(K key, long... ranks) { - return createDissolvingFlux(() -> bloomCommandBuilder.tDigestByRank(key, ranks)); - } - - @Override - public Flux tDigestByRevRank(K key, long... revRanks) { - return createDissolvingFlux(() -> bloomCommandBuilder.tDigestByRevRank(key, revRanks)); - } - - @Override - public Flux tDigestCdf(K key, double... values) { - return createDissolvingFlux(() -> bloomCommandBuilder.tDigestCdf(key, values)); - } - - @Override - public Mono tDigestCreate(K key) { - return createMono(() -> bloomCommandBuilder.tDigestCreate(key)); - } - - @Override - public Mono tDigestCreate(K key, long compression) { - return createMono(() -> bloomCommandBuilder.tDigestCreate(key, compression)); - } - - @Override - public Mono tDigestInfo(K key) { - return createMono(() -> bloomCommandBuilder.tDigestInfo(key)); - } - - @Override - public Mono tDigestMax(K key) { - return createMono(() -> bloomCommandBuilder.tDigestMax(key)); - } - - @Override - public Mono tDigestMerge(K destinationKey, K... sourceKeys) { - return createMono(() -> bloomCommandBuilder.tDigestMerge(destinationKey, sourceKeys)); - } - - @Override - public Mono tDigestMerge(K destinationKey, TDigestMergeOptions options, K... sourceKeys) { - return createMono(() -> bloomCommandBuilder.tDigestMerge(destinationKey, options, sourceKeys)); - } - - @Override - public Mono tDigestMin(K key) { - return createMono(() -> bloomCommandBuilder.tDigestMin(key)); - } - - @Override - public Flux tDigestQuantile(K key, double... quantiles) { - return createDissolvingFlux(() -> bloomCommandBuilder.tDigestQuantile(key, quantiles)); - } - - @Override - public Flux tDigestRank(K key, double... values) { - return createDissolvingFlux(() -> bloomCommandBuilder.tDigestRank(key, values)); - } - - @Override - public Mono tDigestReset(K key) { - return createMono(() -> bloomCommandBuilder.tDigestReset(key)); - } - - @Override - public Flux tDigestRevRank(K key, double... values) { - return createDissolvingFlux(() -> bloomCommandBuilder.tDigestRevRank(key, values)); - } - - @Override - public Mono tDigestTrimmedMean(K key, double lowCutQuantile, double highCutQuantile) { - return createMono(() -> bloomCommandBuilder.tDigestTrimmedMean(key, lowCutQuantile, highCutQuantile)); - } - - @Override - public Mono bfAdd(K key, V item) { - return createMono(() -> bloomCommandBuilder.bfAdd(key, item)); - } - - @Override - public Mono bfCard(K key) { - return createMono(() -> bloomCommandBuilder.bfCard(key)); - } - - @Override - public Mono bfExists(K key, V item) { - return createMono(() -> bloomCommandBuilder.bfExists(key, item)); - } - - @Override - public Mono bfInfo(K key) { - return createMono(() -> bloomCommandBuilder.bfInfo(key)); - } - - @Override - public Mono bfInfo(K key, BloomFilterInfoType infoType) { - return createMono(() -> bloomCommandBuilder.bfInfo(key, infoType)); - } - - @Override - public Flux bfInsert(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.bfInsert(key, items)); - } - - @Override - public Flux bfInsert(K key, BloomFilterInsertOptions options, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.bfInsert(key, options, items)); - } - - @Override - public Flux bfMAdd(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.bfMAdd(key, items)); - } - - @Override - public Flux bfMExists(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.bfMExists(key, items)); - } - - @Override - public Mono bfReserve(K key, double errorRate, long capacity) { - return bfReserve(key, errorRate, capacity, null); - } - - @Override - public Mono bfReserve(K key, double errorRate, long capacity, BloomFilterReserveOptions options) { - return createMono(() -> bloomCommandBuilder.bfReserve(key, errorRate, capacity, options)); - } - - @Override - public Mono cfAdd(K key, V item) { - return createMono(() -> bloomCommandBuilder.cfAdd(key, item)); - } - - @Override - public Mono cfAddNx(K key, V item) { - return createMono(() -> bloomCommandBuilder.cfAddNx(key, item)); - } - - @Override - public Mono cfCount(K key, V item) { - return createMono(() -> bloomCommandBuilder.cfCount(key, item)); - } - - @Override - public Mono cfDel(K key, V item) { - return createMono(() -> bloomCommandBuilder.cfDel(key, item)); - } - - @Override - public Mono cfExists(K key, V item) { - return createMono(() -> bloomCommandBuilder.cfExists(key, item)); - } - - @Override - public Mono cfInfo(K key) { - return createMono(() -> bloomCommandBuilder.cfInfo(key)); - } - - @Override - public Flux cfInsert(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.cfInsert(key, items)); - } - - @Override - public Flux cfInsert(K key, CuckooFilterInsertOptions options, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.cfInsert(key, items, options)); - } - - @Override - public Flux cfInsertNx(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.cfInsertNx(key, items)); - } - - @Override - public Flux cfInsertNx(K key, CuckooFilterInsertOptions options, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.cfInsertNx(key, items, options)); - } - - @Override - public Flux cfMExists(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.cfMExists(key, items)); - } - - @Override - public Mono cfReserve(K key, long capacity) { - return cfReserve(key, capacity, null); - } - - @Override - public Mono cfReserve(K key, long capacity, CuckooFilterReserveOptions options) { - return createMono(() -> bloomCommandBuilder.cfReserve(key, capacity, options)); - } - - @Override - public Mono cmsIncrBy(K key, V item, long increment) { - return createMono(() -> bloomCommandBuilder.cmsIncrBy(key, item, increment)); - } - - @Override - public Flux cmsIncrBy(K key, LongScoredValue... itemIncrements) { - return createDissolvingFlux(() -> bloomCommandBuilder.cmsIncrBy(key, itemIncrements)); - } - - @Override - public Mono cmsInitByProb(K key, double error, double probability) { - return createMono(() -> bloomCommandBuilder.cmsInitByProb(key, error, probability)); - } - - @Override - public Mono cmsInitByDim(K key, long width, long depth) { - return createMono(() -> bloomCommandBuilder.cmsInitByDim(key, width, depth)); - } - - @Override - public Flux cmsQuery(K key, V... items) { - return createDissolvingFlux(() -> bloomCommandBuilder.cmsQuery(key, items)); - } - - @Override - public Mono cmsMerge(K destKey, K... keys) { - return createMono(() -> bloomCommandBuilder.cmsMerge(destKey, keys)); - } - - @Override - public Mono cmsMerge(K destKey, LongScoredValue... sourceKeyWeights) { - return createMono(() -> bloomCommandBuilder.cmsMerge(destKey, sourceKeyWeights)); - } - - @Override - public Mono cmsInfo(K key) { - return createMono(() -> bloomCommandBuilder.cmsInfo(key)); - } + implements RedisModulesReactiveCommands { + + private final StatefulRedisModulesConnection connection; + + private final TimeSeriesCommandBuilder timeSeriesCommandBuilder; + + private final SearchCommandBuilder searchCommandBuilder; + + private final BloomCommandBuilder bloomCommandBuilder; + + public RedisModulesReactiveCommandsImpl(StatefulRedisModulesConnection connection, RedisCodec codec) { + super(connection, codec); + this.connection = connection; + this.timeSeriesCommandBuilder = new TimeSeriesCommandBuilder<>(codec); + this.searchCommandBuilder = new SearchCommandBuilder<>(codec); + this.bloomCommandBuilder = new BloomCommandBuilder<>(codec); + } + + @Override + public StatefulRedisModulesConnection getStatefulConnection() { + return connection; + } + + @Override + public Mono tsCreate(K key, CreateOptions options) { + return createMono(() -> timeSeriesCommandBuilder.create(key, options)); + } + + @Override + public Mono tsAlter(K key, AlterOptions options) { + return createMono(() -> timeSeriesCommandBuilder.alter(key, options)); + } + + @Override + public Mono tsAdd(K key, Sample sample) { + return createMono(() -> timeSeriesCommandBuilder.add(key, sample)); + } + + @Override + public Mono tsAdd(K key, Sample sample, AddOptions options) { + return createMono(() -> timeSeriesCommandBuilder.add(key, sample, options)); + } + + @Override + public Mono tsIncrby(K key, double value) { + return createMono(() -> timeSeriesCommandBuilder.incrby(key, value, null)); + } + + @Override + public Mono tsIncrby(K key, double value, IncrbyOptions options) { + return createMono(() -> timeSeriesCommandBuilder.incrby(key, value, options)); + } + + @Override + public Mono tsDecrby(K key, double value) { + return createMono(() -> timeSeriesCommandBuilder.decrby(key, value, null)); + } + + @Override + public Mono tsDecrby(K key, double value, IncrbyOptions options) { + return createMono(() -> timeSeriesCommandBuilder.decrby(key, value, options)); + } + + @Override + public Flux tsMadd(KeySample... samples) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.madd(samples)); + } + + @Override + public Mono tsCreaterule(K sourceKey, K destKey, CreateRuleOptions options) { + return createMono(() -> timeSeriesCommandBuilder.createRule(sourceKey, destKey, options)); + } + + @Override + public Mono tsDeleterule(K sourceKey, K destKey) { + return createMono(() -> timeSeriesCommandBuilder.deleteRule(sourceKey, destKey)); + } + + @Override + public Flux tsRange(K key, TimeRange range) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.range(key, range)); + } + + @Override + public Flux tsRange(K key, TimeRange range, RangeOptions options) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.range(key, range, options)); + } + + @Override + public Flux tsRevrange(K key, TimeRange range) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.revrange(key, range)); + } + + @Override + public Flux tsRevrange(K key, TimeRange range, RangeOptions options) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.revrange(key, range, options)); + } + + @Override + public Flux> tsMrange(TimeRange range) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrange(range)); + } + + @Override + public Flux> tsMrange(TimeRange range, MRangeOptions options) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrange(range, options)); + } + + @Override + public Flux> tsMrevrange(TimeRange range) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrevrange(range)); + } + + @Override + public Flux> tsMrevrange(TimeRange range, MRangeOptions options) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mrevrange(range, options)); + } + + @Override + public Mono tsGet(K key) { + return createMono(() -> timeSeriesCommandBuilder.get(key)); + } + + @Override + public Flux> tsMget(MGetOptions options) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mget(options)); + } + + @Override + public Flux> tsMget(V... filters) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mget(filters)); + } + + @Override + public Flux> tsMgetWithLabels(V... filters) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.mgetWithLabels(filters)); + } + + @Override + public Flux tsInfo(K key) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.info(key, false)); + } + + @Override + public Flux tsInfoDebug(K key) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.info(key, true)); + } + + @Override + public Flux tsQueryIndex(V... filters) { + return createDissolvingFlux(() -> timeSeriesCommandBuilder.queryIndex(filters)); + } + + @Override + public Mono tsDel(K key, TimeRange timeRange) { + return createMono(() -> timeSeriesCommandBuilder.tsDel(key, timeRange)); + } + + @Override + public Mono ftCreate(K index, Field... fields) { + return ftCreate(index, null, fields); + } + + @Override + public Mono ftCreate(K index, com.redis.lettucemod.search.CreateOptions options, Field... fields) { + return createMono(() -> searchCommandBuilder.create(index, options, fields)); + } + + @Override + public Mono ftDropindex(K index) { + return createMono(() -> searchCommandBuilder.dropIndex(index, false)); + } + + @Override + public Mono ftDropindexDeleteDocs(K index) { + return createMono(() -> searchCommandBuilder.dropIndex(index, true)); + } + + @Override + public Flux ftInfo(K index) { + return createDissolvingFlux(() -> searchCommandBuilder.info(index)); + } + + @Override + public Mono> ftSearch(K index, V query, V... options) { + return createMono(() -> searchCommandBuilder.search(index, query, options)); + } + + @Override + public Mono> ftSearch(K index, V query, SearchOptions options) { + return createMono(() -> searchCommandBuilder.search(index, query, options)); + } + + @Override + public Mono> ftAggregate(K index, V query, V... options) { + return createMono(() -> searchCommandBuilder.aggregate(index, query, options)); + } + + @Override + public Mono> ftAggregate(K index, V query, AggregateOptions options) { + return createMono(() -> searchCommandBuilder.aggregate(index, query, options)); + } + + @Override + public Mono> ftAggregate(K index, V query, CursorOptions cursor) { + return createMono(() -> searchCommandBuilder.aggregate(index, query, cursor, null)); + } + + @Override + public Mono> ftAggregate(K index, V query, CursorOptions cursor, + AggregateOptions options) { + return createMono(() -> searchCommandBuilder.aggregate(index, query, cursor, options)); + } + + @Override + public Mono> ftCursorRead(K index, long cursor) { + return createMono(() -> searchCommandBuilder.cursorRead(index, cursor, null)); + } + + @Override + public Mono> ftCursorRead(K index, long cursor, long count) { + return createMono(() -> searchCommandBuilder.cursorRead(index, cursor, count)); + } + + @Override + public Mono ftCursorDelete(K index, long cursor) { + return createMono(() -> searchCommandBuilder.cursorDelete(index, cursor)); + } + + @Override + public Mono ftSugadd(K key, Suggestion suggestion) { + return createMono(() -> searchCommandBuilder.sugadd(key, suggestion)); + } + + @Override + public Mono ftSugaddIncr(K key, Suggestion suggestion) { + return createMono(() -> searchCommandBuilder.sugaddIncr(key, suggestion)); + } + + @Override + public Flux> ftSugget(K key, V prefix) { + return createDissolvingFlux(() -> searchCommandBuilder.sugget(key, prefix)); + } + + @Override + public Flux> ftSugget(K key, V prefix, SuggetOptions options) { + return createDissolvingFlux(() -> searchCommandBuilder.sugget(key, prefix, options)); + } + + @Override + public Mono ftSugdel(K key, V string) { + return createMono(() -> searchCommandBuilder.sugdel(key, string)); + } + + @Override + public Mono ftSuglen(K key) { + return createMono(() -> searchCommandBuilder.suglen(key)); + } + + @Override + public Mono ftAlter(K index, Field field) { + return createMono(() -> searchCommandBuilder.alter(index, field)); + } + + @Override + public Mono ftAliasadd(K name, K index) { + return createMono(() -> searchCommandBuilder.aliasAdd(name, index)); + } + + @Override + public Mono ftAliasupdate(K name, K index) { + return createMono(() -> searchCommandBuilder.aliasUpdate(name, index)); + } + + @Override + public Mono ftAliasdel(K name) { + return createMono(() -> searchCommandBuilder.aliasDel(name)); + } + + @Override + public Flux ftList() { + return createDissolvingFlux(searchCommandBuilder::list); + } + + @Override + public Flux ftTagvals(K index, K field) { + return createDissolvingFlux(() -> searchCommandBuilder.tagVals(index, field)); + } + + @Override + public Mono ftDictadd(K dict, V... terms) { + return createMono(() -> searchCommandBuilder.dictadd(dict, terms)); + } + + @Override + public Mono ftDictdel(K dict, V... terms) { + return createMono(() -> searchCommandBuilder.dictdel(dict, terms)); + } + + @Override + public Flux ftDictdump(K dict) { + return createDissolvingFlux(() -> searchCommandBuilder.dictdump(dict)); + } + + @Override + public Flux> topKAdd(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.topKAdd(key, items)); + } + + @Override + public Flux> topKIncrBy(K key, LongScoredValue... itemIncrements) { + return createDissolvingFlux(() -> bloomCommandBuilder.topKIncrBy(key, itemIncrements)); + } + + @Override + public Mono topKInfo(K key) { + return createMono(() -> bloomCommandBuilder.topKInfo(key)); + } + + @Override + public Flux topKList(K key) { + return createDissolvingFlux(() -> bloomCommandBuilder.topKList(key)); + } + + @Override + public Flux> topKListWithScores(K key) { + return createDissolvingFlux(() -> bloomCommandBuilder.topKListWithScores(key)); + } + + @Override + public Flux topKQuery(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.topKQuery(key, items)); + } + + @Override + public Mono topKReserve(K key, long k) { + return createMono(() -> bloomCommandBuilder.topKReserve(key, k)); + } + + @Override + public Mono topKReserve(K key, long k, long width, long depth, double decay) { + return createMono(() -> bloomCommandBuilder.topKReserve(key, k, width, depth, decay)); + } + + @Override + public Mono tDigestAdd(K key, double... values) { + return createMono(() -> bloomCommandBuilder.tDigestAdd(key, values)); + } + + @Override + public Flux tDigestByRank(K key, long... ranks) { + return createDissolvingFlux(() -> bloomCommandBuilder.tDigestByRank(key, ranks)); + } + + @Override + public Flux tDigestByRevRank(K key, long... revRanks) { + return createDissolvingFlux(() -> bloomCommandBuilder.tDigestByRevRank(key, revRanks)); + } + + @Override + public Flux tDigestCdf(K key, double... values) { + return createDissolvingFlux(() -> bloomCommandBuilder.tDigestCdf(key, values)); + } + + @Override + public Mono tDigestCreate(K key) { + return createMono(() -> bloomCommandBuilder.tDigestCreate(key)); + } + + @Override + public Mono tDigestCreate(K key, long compression) { + return createMono(() -> bloomCommandBuilder.tDigestCreate(key, compression)); + } + + @Override + public Mono tDigestInfo(K key) { + return createMono(() -> bloomCommandBuilder.tDigestInfo(key)); + } + + @Override + public Mono tDigestMax(K key) { + return createMono(() -> bloomCommandBuilder.tDigestMax(key)); + } + + @Override + public Mono tDigestMerge(K destinationKey, K... sourceKeys) { + return createMono(() -> bloomCommandBuilder.tDigestMerge(destinationKey, sourceKeys)); + } + + @Override + public Mono tDigestMerge(K destinationKey, TDigestMergeOptions options, K... sourceKeys) { + return createMono(() -> bloomCommandBuilder.tDigestMerge(destinationKey, options, sourceKeys)); + } + + @Override + public Mono tDigestMin(K key) { + return createMono(() -> bloomCommandBuilder.tDigestMin(key)); + } + + @Override + public Flux tDigestQuantile(K key, double... quantiles) { + return createDissolvingFlux(() -> bloomCommandBuilder.tDigestQuantile(key, quantiles)); + } + + @Override + public Flux tDigestRank(K key, double... values) { + return createDissolvingFlux(() -> bloomCommandBuilder.tDigestRank(key, values)); + } + + @Override + public Mono tDigestReset(K key) { + return createMono(() -> bloomCommandBuilder.tDigestReset(key)); + } + + @Override + public Flux tDigestRevRank(K key, double... values) { + return createDissolvingFlux(() -> bloomCommandBuilder.tDigestRevRank(key, values)); + } + + @Override + public Mono tDigestTrimmedMean(K key, double lowCutQuantile, double highCutQuantile) { + return createMono(() -> bloomCommandBuilder.tDigestTrimmedMean(key, lowCutQuantile, highCutQuantile)); + } + + @Override + public Mono bfAdd(K key, V item) { + return createMono(() -> bloomCommandBuilder.bfAdd(key, item)); + } + + @Override + public Mono bfCard(K key) { + return createMono(() -> bloomCommandBuilder.bfCard(key)); + } + + @Override + public Mono bfExists(K key, V item) { + return createMono(() -> bloomCommandBuilder.bfExists(key, item)); + } + + @Override + public Mono bfInfo(K key) { + return createMono(() -> bloomCommandBuilder.bfInfo(key)); + } + + @Override + public Mono bfInfo(K key, BloomFilterInfoType infoType) { + return createMono(() -> bloomCommandBuilder.bfInfo(key, infoType)); + } + + @Override + public Flux bfInsert(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.bfInsert(key, items)); + } + + @Override + public Flux bfInsert(K key, BloomFilterInsertOptions options, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.bfInsert(key, options, items)); + } + + @Override + public Flux bfMAdd(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.bfMAdd(key, items)); + } + + @Override + public Flux bfMExists(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.bfMExists(key, items)); + } + + @Override + public Mono bfReserve(K key, double errorRate, long capacity) { + return bfReserve(key, errorRate, capacity, null); + } + + @Override + public Mono bfReserve(K key, double errorRate, long capacity, BloomFilterReserveOptions options) { + return createMono(() -> bloomCommandBuilder.bfReserve(key, errorRate, capacity, options)); + } + + @Override + public Mono cfAdd(K key, V item) { + return createMono(() -> bloomCommandBuilder.cfAdd(key, item)); + } + + @Override + public Mono cfAddNx(K key, V item) { + return createMono(() -> bloomCommandBuilder.cfAddNx(key, item)); + } + + @Override + public Mono cfCount(K key, V item) { + return createMono(() -> bloomCommandBuilder.cfCount(key, item)); + } + + @Override + public Mono cfDel(K key, V item) { + return createMono(() -> bloomCommandBuilder.cfDel(key, item)); + } + + @Override + public Mono cfExists(K key, V item) { + return createMono(() -> bloomCommandBuilder.cfExists(key, item)); + } + + @Override + public Mono cfInfo(K key) { + return createMono(() -> bloomCommandBuilder.cfInfo(key)); + } + + @Override + public Flux cfInsert(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.cfInsert(key, items)); + } + + @Override + public Flux cfInsert(K key, CuckooFilterInsertOptions options, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.cfInsert(key, items, options)); + } + + @Override + public Flux cfInsertNx(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.cfInsertNx(key, items)); + } + + @Override + public Flux cfInsertNx(K key, CuckooFilterInsertOptions options, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.cfInsertNx(key, items, options)); + } + + @Override + public Flux cfMExists(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.cfMExists(key, items)); + } + + @Override + public Mono cfReserve(K key, long capacity) { + return cfReserve(key, capacity, null); + } + + @Override + public Mono cfReserve(K key, long capacity, CuckooFilterReserveOptions options) { + return createMono(() -> bloomCommandBuilder.cfReserve(key, capacity, options)); + } + + @Override + public Mono cmsIncrBy(K key, V item, long increment) { + return createMono(() -> bloomCommandBuilder.cmsIncrBy(key, item, increment)); + } + + @Override + public Flux cmsIncrBy(K key, LongScoredValue... itemIncrements) { + return createDissolvingFlux(() -> bloomCommandBuilder.cmsIncrBy(key, itemIncrements)); + } + + @Override + public Mono cmsInitByProb(K key, double error, double probability) { + return createMono(() -> bloomCommandBuilder.cmsInitByProb(key, error, probability)); + } + + @Override + public Mono cmsInitByDim(K key, long width, long depth) { + return createMono(() -> bloomCommandBuilder.cmsInitByDim(key, width, depth)); + } + + @Override + public Flux cmsQuery(K key, V... items) { + return createDissolvingFlux(() -> bloomCommandBuilder.cmsQuery(key, items)); + } + + @Override + public Mono cmsMerge(K destKey, K... keys) { + return createMono(() -> bloomCommandBuilder.cmsMerge(destKey, keys)); + } + + @Override + public Mono cmsMerge(K destKey, LongScoredValue... sourceKeyWeights) { + return createMono(() -> bloomCommandBuilder.cmsMerge(destKey, sourceKeyWeights)); + } + + @Override + public Mono cmsInfo(K key) { + return createMono(() -> bloomCommandBuilder.cmsInfo(key)); + } + } diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/api/async/RediSearchAsyncCommands.java b/core/lettucemod/src/main/java/com/redis/lettucemod/api/async/RediSearchAsyncCommands.java index b25fe07..6304341 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/api/async/RediSearchAsyncCommands.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/api/async/RediSearchAsyncCommands.java @@ -30,11 +30,11 @@ public interface RediSearchAsyncCommands { RedisFuture> ftList(); - RedisFuture> ftSearch(K index, V query); + RedisFuture> ftSearch(K index, V query, V... options); RedisFuture> ftSearch(K index, V query, SearchOptions options); - RedisFuture> ftAggregate(K index, V query); + RedisFuture> ftAggregate(K index, V query, V... options); RedisFuture> ftAggregate(K index, V query, AggregateOptions options); diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/api/reactive/RediSearchReactiveCommands.java b/core/lettucemod/src/main/java/com/redis/lettucemod/api/reactive/RediSearchReactiveCommands.java index 0bcbb9d..a0c624c 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/api/reactive/RediSearchReactiveCommands.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/api/reactive/RediSearchReactiveCommands.java @@ -38,11 +38,11 @@ public interface RediSearchReactiveCommands { Flux ftList(); - Mono> ftSearch(K index, V query); + Mono> ftSearch(K index, V query, V... options); Mono> ftSearch(K index, V query, SearchOptions options); - Mono> ftAggregate(K index, V query); + Mono> ftAggregate(K index, V query, V... options); Mono> ftAggregate(K index, V query, AggregateOptions options); diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/api/sync/RediSearchCommands.java b/core/lettucemod/src/main/java/com/redis/lettucemod/api/sync/RediSearchCommands.java index 2365fc5..066a564 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/api/sync/RediSearchCommands.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/api/sync/RediSearchCommands.java @@ -32,11 +32,11 @@ public interface RediSearchCommands { */ List ftList(); - SearchResults ftSearch(K index, V query); + SearchResults ftSearch(K index, V query, V... options); SearchResults ftSearch(K index, V query, SearchOptions options); - AggregateResults ftAggregate(K index, V query); + AggregateResults ftAggregate(K index, V query, V... options); AggregateResults ftAggregate(K index, V query, AggregateOptions options); diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterAsyncCommandsImpl.java b/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterAsyncCommandsImpl.java index 84e8ca4..84bed95 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterAsyncCommandsImpl.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterAsyncCommandsImpl.java @@ -134,8 +134,8 @@ public RedisFuture> ftList() { } @Override - public RedisFuture> ftSearch(K index, V query) { - return delegate.ftSearch(index, query); + public RedisFuture> ftSearch(K index, V query, V... options) { + return delegate.ftSearch(index, query, options); } @Override @@ -144,8 +144,8 @@ public RedisFuture> ftSearch(K index, V query, SearchOptions } @Override - public RedisFuture> ftAggregate(K index, V query) { - return delegate.ftAggregate(index, query); + public RedisFuture> ftAggregate(K index, V query, V... options) { + return delegate.ftAggregate(index, query, options); } @Override diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterReactiveCommandsImpl.java b/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterReactiveCommandsImpl.java index d3b5802..11ca738 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterReactiveCommandsImpl.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/cluster/RedisModulesAdvancedClusterReactiveCommandsImpl.java @@ -51,608 +51,609 @@ import reactor.core.publisher.Mono; @SuppressWarnings("unchecked") -public class RedisModulesAdvancedClusterReactiveCommandsImpl extends - RedisAdvancedClusterReactiveCommandsImpl implements RedisModulesAdvancedClusterReactiveCommands { - - private final RedisModulesReactiveCommandsImpl delegate; - - public RedisModulesAdvancedClusterReactiveCommandsImpl(StatefulRedisModulesClusterConnection connection, - RedisCodec codec) { - super(connection, codec); - this.delegate = new RedisModulesReactiveCommandsImpl<>(connection, codec); - } - - @Override - public RedisModulesAdvancedClusterReactiveCommands getConnection(String nodeId) { - return (RedisModulesAdvancedClusterReactiveCommands) super.getConnection(nodeId); - } - - @Override - public RedisModulesAdvancedClusterReactiveCommands getConnection(String host, int port) { - return (RedisModulesAdvancedClusterReactiveCommands) super.getConnection(host, port); - } - - @Override - public StatefulRedisModulesClusterConnection getStatefulConnection() { - return (StatefulRedisModulesClusterConnection) super.getStatefulConnection(); - } - - @Override - public Mono ftCreate(K index, Field... fields) { - return ftCreate(index, null, fields); - } - - @Override - public Mono ftCreate(K index, com.redis.lettucemod.search.CreateOptions options, Field... fields) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftCreate(index, options, fields)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Mono ftDropindex(K index) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftDropindex(index)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Mono ftDropindexDeleteDocs(K index) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftDropindexDeleteDocs(index)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Mono ftAlter(K index, Field field) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftAlter(index, field)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Flux ftInfo(K index) { - return delegate.ftInfo(index); - } - - @Override - public Mono ftAliasadd(K name, K index) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftAliasadd(name, index)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Mono ftAliasupdate(K name, K index) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftAliasupdate(name, index)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Mono ftAliasdel(K name) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftAliasdel(name)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Flux ftList() { - return delegate.ftList(); - } - - @Override - public Mono> ftSearch(K index, V query) { - return delegate.ftSearch(index, query); - } - - @Override - public Mono> ftSearch(K index, V query, SearchOptions options) { - return delegate.ftSearch(index, query, options); - } - - @Override - public Mono> ftAggregate(K index, V query) { - return delegate.ftAggregate(index, query); - } - - @Override - public Mono> ftAggregate(K index, V query, AggregateOptions options) { - return delegate.ftAggregate(index, query, options); - } - - @Override - public Mono> ftAggregate(K index, V query, CursorOptions cursor) { - return delegate.ftAggregate(index, query, cursor); - } - - @Override - public Mono> ftAggregate(K index, V query, CursorOptions cursor, - AggregateOptions options) { - return delegate.ftAggregate(index, query, cursor, options); - } - - @Override - public Mono> ftCursorRead(K index, long cursor) { - return delegate.ftCursorRead(index, cursor); - } - - @Override - public Mono> ftCursorRead(K index, long cursor, long count) { - return delegate.ftCursorRead(index, cursor, count); - } - - @Override - public Mono ftCursorDelete(K index, long cursor) { - return delegate.ftCursorDelete(index, cursor); - } - - @Override - public Flux ftTagvals(K index, K field) { - return delegate.ftTagvals(index, field); - } - - @Override - public Mono ftSugadd(K key, Suggestion suggestion) { - return delegate.ftSugadd(key, suggestion); - } - - @Override - public Mono ftSugaddIncr(K key, Suggestion suggestion) { - return delegate.ftSugaddIncr(key, suggestion); - } - - @Override - public Flux> ftSugget(K key, V prefix) { - return delegate.ftSugget(key, prefix); - } - - @Override - public Flux> ftSugget(K key, V prefix, SuggetOptions options) { - return delegate.ftSugget(key, prefix, options); - } - - @Override - public Mono ftSugdel(K key, V string) { - return delegate.ftSugdel(key, string); - } - - @Override - public Mono ftSuglen(K key) { - return delegate.ftSuglen(key); - } - - @Override - public Mono ftDictadd(K dict, V... terms) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftDictadd(dict, terms)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Mono ftDictdel(K dict, V... terms) { - Map> publishers = executeOnUpstream( - commands -> ((RedisModulesReactiveCommands) commands).ftDictdel(dict, terms)); - return Flux.merge(publishers.values()).last(); - } - - @Override - public Flux ftDictdump(K dict) { - return delegate.ftDictdump(dict); - } - - @Override - public Flux> topKAdd(K key, V... items) { - return delegate.topKAdd(key, items); - } - - @Override - public Flux> topKIncrBy(K key, LongScoredValue... itemIncrements) { - return delegate.topKIncrBy(key, itemIncrements); - } - - @Override - public Mono topKInfo(K key) { - return delegate.topKInfo(key); - } - - @Override - public Flux topKList(K key) { - return delegate.topKList(key); - } - - @Override - public Flux> topKListWithScores(K key) { - return delegate.topKListWithScores(key); - } - - @Override - public Flux topKQuery(K key, V... items) { - return delegate.topKQuery(key, items); - } - - @Override - public Mono topKReserve(K key, long k) { - return delegate.topKReserve(key, k); - } - - @Override - public Mono topKReserve(K key, long k, long width, long depth, double decay) { - return delegate.topKReserve(key, k, width, depth, decay); - } - - @Override - public Mono tDigestAdd(K key, double... value) { - return delegate.tDigestAdd(key, value); - } - - @Override - public Flux tDigestByRank(K key, long... ranks) { - return delegate.tDigestByRank(key, ranks); - } - - @Override - public Flux tDigestByRevRank(K key, long... revRanks) { - return delegate.tDigestByRevRank(key, revRanks); - } - - @Override - public Flux tDigestCdf(K key, double... values) { - return delegate.tDigestCdf(key, values); - } - - @Override - public Mono tDigestCreate(K key) { - return delegate.tDigestCreate(key); - } - - @Override - public Mono tDigestCreate(K key, long compression) { - return delegate.tDigestCreate(key, compression); - } - - @Override - public Mono tDigestInfo(K key) { - return delegate.tDigestInfo(key); - } - - @Override - public Mono tDigestMax(K key) { - return delegate.tDigestMax(key); - } - - @Override - public Mono tDigestMerge(K destinationKey, K... sourceKeys) { - return delegate.tDigestMerge(destinationKey, sourceKeys); - } - - @Override - public Mono tDigestMerge(K destinationKey, TDigestMergeOptions options, K... sourceKeys) { - return delegate.tDigestMerge(destinationKey, options, sourceKeys); - } - - @Override - public Mono tDigestMin(K key) { - return delegate.tDigestMin(key); - } - - @Override - public Flux tDigestQuantile(K key, double... quantiles) { - return delegate.tDigestQuantile(key, quantiles); - } - - @Override - public Flux tDigestRank(K key, double... values) { - return delegate.tDigestRank(key, values); - } - - @Override - public Mono tDigestReset(K key) { - return delegate.tDigestReset(key); - } - - @Override - public Flux tDigestRevRank(K key, double... values) { - return delegate.tDigestRevRank(key, values); - } - - @Override - public Mono tDigestTrimmedMean(K key, double lowCutQuantile, double highCutQuantile) { - return delegate.tDigestTrimmedMean(key, lowCutQuantile, highCutQuantile); - } - - @Override - public Mono tsCreate(K key, CreateOptions options) { - return delegate.tsCreate(key, options); - } - - @Override - public Mono tsAlter(K key, AlterOptions options) { - return delegate.tsAlter(key, options); - } - - @Override - public Mono tsAdd(K key, Sample sample) { - return delegate.tsAdd(key, sample); - } - - @Override - public Mono tsAdd(K key, Sample sample, AddOptions options) { - return delegate.tsAdd(key, sample, options); - } - - @Override - public Flux tsMadd(KeySample... samples) { - return delegate.tsMadd(samples); - } - - @Override - public Mono tsDecrby(K key, double value) { - return delegate.tsDecrby(key, value); - } - - @Override - public Mono tsDecrby(K key, double value, IncrbyOptions options) { - return delegate.tsDecrby(key, value, options); - } - - @Override - public Mono tsIncrby(K key, double value) { - return delegate.tsIncrby(key, value); - } - - @Override - public Mono tsIncrby(K key, double value, IncrbyOptions options) { - return delegate.tsIncrby(key, value, options); - } - - @Override - public Mono tsCreaterule(K sourceKey, K destKey, CreateRuleOptions options) { - return delegate.tsCreaterule(sourceKey, destKey, options); - } - - @Override - public Mono tsDeleterule(K sourceKey, K destKey) { - return delegate.tsDeleterule(sourceKey, destKey); - } - - @Override - public Flux tsRange(K key, TimeRange range) { - return delegate.tsRange(key, range); - } - - @Override - public Flux tsRange(K key, TimeRange range, RangeOptions options) { - return delegate.tsRange(key, range, options); - } - - @Override - public Flux tsRevrange(K key, TimeRange range) { - return delegate.tsRevrange(key, range); - } - - @Override - public Flux tsRevrange(K key, TimeRange range, RangeOptions options) { - return delegate.tsRevrange(key, range, options); - } - - @Override - public Flux> tsMrange(TimeRange range) { - return delegate.tsMrange(range); - } - - @Override - public Flux> tsMrange(TimeRange range, MRangeOptions options) { - return delegate.tsMrange(range, options); - } - - @Override - public Flux> tsMrevrange(TimeRange range) { - return delegate.tsMrevrange(range); - } - - @Override - public Flux> tsMrevrange(TimeRange range, MRangeOptions options) { - return delegate.tsMrevrange(range, options); - } - - @Override - public Mono tsGet(K key) { - return delegate.tsGet(key); - } - - @Override - public Flux> tsMget(MGetOptions options) { - return delegate.tsMget(options); - } - - @Override - public Flux> tsMget(V... filters) { - return delegate.tsMget(filters); - } - - @Override - public Flux> tsMgetWithLabels(V... filters) { - return delegate.tsMgetWithLabels(filters); - } - - @Override - public Flux tsInfo(K key) { - return delegate.tsInfo(key); - } - - @Override - public Flux tsInfoDebug(K key) { - return delegate.tsInfoDebug(key); - } - - @Override - public Flux tsQueryIndex(V... filters) { - return delegate.tsQueryIndex(filters); - } - - @Override - public Mono tsDel(K key, TimeRange timeRange) { - return delegate.tsDel(key, timeRange); - } - - @Override - public Mono bfAdd(K key, V item) { - return delegate.bfAdd(key, item); - } - - @Override - public Mono bfCard(K key) { - return delegate.bfCard(key); - } - - @Override - public Mono bfExists(K key, V item) { - return delegate.bfExists(key, item); - } - - @Override - public Mono bfInfo(K key) { - return delegate.bfInfo(key); - } - - @Override - public Mono bfInfo(K key, BloomFilterInfoType infoType) { - return delegate.bfInfo(key, infoType); - } - - @Override - public Flux bfInsert(K key, V... items) { - return delegate.bfInsert(key, items); - } - - @Override - public Flux bfInsert(K key, BloomFilterInsertOptions options, V... items) { - return delegate.bfInsert(key, options, items); - } - - @Override - public Flux bfMAdd(K key, V... items) { - return delegate.bfMAdd(key, items); - } - - @Override - public Flux bfMExists(K key, V... items) { - return delegate.bfMExists(key, items); - } - - @Override - public Mono bfReserve(K key, double errorRate, long capacity) { - return delegate.bfReserve(key, errorRate, capacity); - } - - @Override - public Mono bfReserve(K key, double errorRate, long capacity, BloomFilterReserveOptions options) { - return delegate.bfReserve(key, errorRate, capacity, options); - } - - @Override - public Mono cfAdd(K key, V item) { - return delegate.cfAdd(key, item); - } - - @Override - public Mono cfAddNx(K key, V item) { - return delegate.cfAddNx(key, item); - } - - @Override - public Mono cfCount(K key, V item) { - return delegate.cfCount(key, item); - } - - @Override - public Mono cfDel(K key, V item) { - return delegate.cfDel(key, item); - } - - @Override - public Mono cfExists(K key, V item) { - return delegate.cfExists(key, item); - } - - @Override - public Mono cfInfo(K key) { - return delegate.cfInfo(key); - } - - @Override - public Flux cfInsert(K key, V... items) { - return delegate.cfInsert(key, items); - } - - @Override - public Flux cfInsert(K key, CuckooFilterInsertOptions options, V... items) { - return delegate.cfInsert(key, options, items); - } - - @Override - public Flux cfInsertNx(K key, V... items) { - return delegate.cfInsertNx(key, items); - } - - @Override - public Flux cfInsertNx(K key, CuckooFilterInsertOptions options, V... items) { - return delegate.cfInsertNx(key, options, items); - } - - @Override - public Flux cfMExists(K key, V... items) { - return delegate.cfMExists(key, items); - } - - @Override - public Mono cfReserve(K key, long capacity) { - return delegate.cfReserve(key, capacity); - } - - @Override - public Mono cfReserve(K key, long capacity, CuckooFilterReserveOptions options) { - return delegate.cfReserve(key, capacity, options); - } - - @Override - public Mono cmsIncrBy(K key, V item, long increment) { - return delegate.cmsIncrBy(key, item, increment); - } - - @Override - public Flux cmsIncrBy(K key, LongScoredValue... itemIncrements) { - return delegate.cmsIncrBy(key, itemIncrements); - } - - @Override - public Mono cmsInitByProb(K key, double error, double probability) { - return delegate.cmsInitByProb(key, error, probability); - } - - @Override - public Mono cmsInitByDim(K key, long width, long depth) { - return delegate.cmsInitByDim(key, width, depth); - } - - @Override - public Flux cmsQuery(K key, V... items) { - return delegate.cmsQuery(key, items); - } - - @Override - public Mono cmsMerge(K destKey, K... keys) { - return delegate.cmsMerge(destKey, keys); - } - - @Override - public Mono cmsMerge(K destKey, LongScoredValue... sourceKeyWeights) { - return delegate.cmsMerge(destKey, sourceKeyWeights); - } - - @Override - public Mono cmsInfo(K key) { - return delegate.cmsInfo(key); - } +public class RedisModulesAdvancedClusterReactiveCommandsImpl extends RedisAdvancedClusterReactiveCommandsImpl + implements RedisModulesAdvancedClusterReactiveCommands { + + private final RedisModulesReactiveCommandsImpl delegate; + + public RedisModulesAdvancedClusterReactiveCommandsImpl(StatefulRedisModulesClusterConnection connection, + RedisCodec codec) { + super(connection, codec); + this.delegate = new RedisModulesReactiveCommandsImpl<>(connection, codec); + } + + @Override + public RedisModulesAdvancedClusterReactiveCommands getConnection(String nodeId) { + return (RedisModulesAdvancedClusterReactiveCommands) super.getConnection(nodeId); + } + + @Override + public RedisModulesAdvancedClusterReactiveCommands getConnection(String host, int port) { + return (RedisModulesAdvancedClusterReactiveCommands) super.getConnection(host, port); + } + + @Override + public StatefulRedisModulesClusterConnection getStatefulConnection() { + return (StatefulRedisModulesClusterConnection) super.getStatefulConnection(); + } + + @Override + public Mono ftCreate(K index, Field... fields) { + return ftCreate(index, null, fields); + } + + @Override + public Mono ftCreate(K index, com.redis.lettucemod.search.CreateOptions options, Field... fields) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftCreate(index, options, fields)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Mono ftDropindex(K index) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftDropindex(index)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Mono ftDropindexDeleteDocs(K index) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftDropindexDeleteDocs(index)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Mono ftAlter(K index, Field field) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftAlter(index, field)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Flux ftInfo(K index) { + return delegate.ftInfo(index); + } + + @Override + public Mono ftAliasadd(K name, K index) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftAliasadd(name, index)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Mono ftAliasupdate(K name, K index) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftAliasupdate(name, index)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Mono ftAliasdel(K name) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftAliasdel(name)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Flux ftList() { + return delegate.ftList(); + } + + @Override + public Mono> ftSearch(K index, V query, V... options) { + return delegate.ftSearch(index, query, options); + } + + @Override + public Mono> ftSearch(K index, V query, SearchOptions options) { + return delegate.ftSearch(index, query, options); + } + + @Override + public Mono> ftAggregate(K index, V query, V... options) { + return delegate.ftAggregate(index, query, options); + } + + @Override + public Mono> ftAggregate(K index, V query, AggregateOptions options) { + return delegate.ftAggregate(index, query, options); + } + + @Override + public Mono> ftAggregate(K index, V query, CursorOptions cursor) { + return delegate.ftAggregate(index, query, cursor); + } + + @Override + public Mono> ftAggregate(K index, V query, CursorOptions cursor, + AggregateOptions options) { + return delegate.ftAggregate(index, query, cursor, options); + } + + @Override + public Mono> ftCursorRead(K index, long cursor) { + return delegate.ftCursorRead(index, cursor); + } + + @Override + public Mono> ftCursorRead(K index, long cursor, long count) { + return delegate.ftCursorRead(index, cursor, count); + } + + @Override + public Mono ftCursorDelete(K index, long cursor) { + return delegate.ftCursorDelete(index, cursor); + } + + @Override + public Flux ftTagvals(K index, K field) { + return delegate.ftTagvals(index, field); + } + + @Override + public Mono ftSugadd(K key, Suggestion suggestion) { + return delegate.ftSugadd(key, suggestion); + } + + @Override + public Mono ftSugaddIncr(K key, Suggestion suggestion) { + return delegate.ftSugaddIncr(key, suggestion); + } + + @Override + public Flux> ftSugget(K key, V prefix) { + return delegate.ftSugget(key, prefix); + } + + @Override + public Flux> ftSugget(K key, V prefix, SuggetOptions options) { + return delegate.ftSugget(key, prefix, options); + } + + @Override + public Mono ftSugdel(K key, V string) { + return delegate.ftSugdel(key, string); + } + + @Override + public Mono ftSuglen(K key) { + return delegate.ftSuglen(key); + } + + @Override + public Mono ftDictadd(K dict, V... terms) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftDictadd(dict, terms)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Mono ftDictdel(K dict, V... terms) { + Map> publishers = executeOnUpstream( + commands -> ((RedisModulesReactiveCommands) commands).ftDictdel(dict, terms)); + return Flux.merge(publishers.values()).last(); + } + + @Override + public Flux ftDictdump(K dict) { + return delegate.ftDictdump(dict); + } + + @Override + public Flux> topKAdd(K key, V... items) { + return delegate.topKAdd(key, items); + } + + @Override + public Flux> topKIncrBy(K key, LongScoredValue... itemIncrements) { + return delegate.topKIncrBy(key, itemIncrements); + } + + @Override + public Mono topKInfo(K key) { + return delegate.topKInfo(key); + } + + @Override + public Flux topKList(K key) { + return delegate.topKList(key); + } + + @Override + public Flux> topKListWithScores(K key) { + return delegate.topKListWithScores(key); + } + + @Override + public Flux topKQuery(K key, V... items) { + return delegate.topKQuery(key, items); + } + + @Override + public Mono topKReserve(K key, long k) { + return delegate.topKReserve(key, k); + } + + @Override + public Mono topKReserve(K key, long k, long width, long depth, double decay) { + return delegate.topKReserve(key, k, width, depth, decay); + } + + @Override + public Mono tDigestAdd(K key, double... value) { + return delegate.tDigestAdd(key, value); + } + + @Override + public Flux tDigestByRank(K key, long... ranks) { + return delegate.tDigestByRank(key, ranks); + } + + @Override + public Flux tDigestByRevRank(K key, long... revRanks) { + return delegate.tDigestByRevRank(key, revRanks); + } + + @Override + public Flux tDigestCdf(K key, double... values) { + return delegate.tDigestCdf(key, values); + } + + @Override + public Mono tDigestCreate(K key) { + return delegate.tDigestCreate(key); + } + + @Override + public Mono tDigestCreate(K key, long compression) { + return delegate.tDigestCreate(key, compression); + } + + @Override + public Mono tDigestInfo(K key) { + return delegate.tDigestInfo(key); + } + + @Override + public Mono tDigestMax(K key) { + return delegate.tDigestMax(key); + } + + @Override + public Mono tDigestMerge(K destinationKey, K... sourceKeys) { + return delegate.tDigestMerge(destinationKey, sourceKeys); + } + + @Override + public Mono tDigestMerge(K destinationKey, TDigestMergeOptions options, K... sourceKeys) { + return delegate.tDigestMerge(destinationKey, options, sourceKeys); + } + + @Override + public Mono tDigestMin(K key) { + return delegate.tDigestMin(key); + } + + @Override + public Flux tDigestQuantile(K key, double... quantiles) { + return delegate.tDigestQuantile(key, quantiles); + } + + @Override + public Flux tDigestRank(K key, double... values) { + return delegate.tDigestRank(key, values); + } + + @Override + public Mono tDigestReset(K key) { + return delegate.tDigestReset(key); + } + + @Override + public Flux tDigestRevRank(K key, double... values) { + return delegate.tDigestRevRank(key, values); + } + + @Override + public Mono tDigestTrimmedMean(K key, double lowCutQuantile, double highCutQuantile) { + return delegate.tDigestTrimmedMean(key, lowCutQuantile, highCutQuantile); + } + + @Override + public Mono tsCreate(K key, CreateOptions options) { + return delegate.tsCreate(key, options); + } + + @Override + public Mono tsAlter(K key, AlterOptions options) { + return delegate.tsAlter(key, options); + } + + @Override + public Mono tsAdd(K key, Sample sample) { + return delegate.tsAdd(key, sample); + } + + @Override + public Mono tsAdd(K key, Sample sample, AddOptions options) { + return delegate.tsAdd(key, sample, options); + } + + @Override + public Flux tsMadd(KeySample... samples) { + return delegate.tsMadd(samples); + } + + @Override + public Mono tsDecrby(K key, double value) { + return delegate.tsDecrby(key, value); + } + + @Override + public Mono tsDecrby(K key, double value, IncrbyOptions options) { + return delegate.tsDecrby(key, value, options); + } + + @Override + public Mono tsIncrby(K key, double value) { + return delegate.tsIncrby(key, value); + } + + @Override + public Mono tsIncrby(K key, double value, IncrbyOptions options) { + return delegate.tsIncrby(key, value, options); + } + + @Override + public Mono tsCreaterule(K sourceKey, K destKey, CreateRuleOptions options) { + return delegate.tsCreaterule(sourceKey, destKey, options); + } + + @Override + public Mono tsDeleterule(K sourceKey, K destKey) { + return delegate.tsDeleterule(sourceKey, destKey); + } + + @Override + public Flux tsRange(K key, TimeRange range) { + return delegate.tsRange(key, range); + } + + @Override + public Flux tsRange(K key, TimeRange range, RangeOptions options) { + return delegate.tsRange(key, range, options); + } + + @Override + public Flux tsRevrange(K key, TimeRange range) { + return delegate.tsRevrange(key, range); + } + + @Override + public Flux tsRevrange(K key, TimeRange range, RangeOptions options) { + return delegate.tsRevrange(key, range, options); + } + + @Override + public Flux> tsMrange(TimeRange range) { + return delegate.tsMrange(range); + } + + @Override + public Flux> tsMrange(TimeRange range, MRangeOptions options) { + return delegate.tsMrange(range, options); + } + + @Override + public Flux> tsMrevrange(TimeRange range) { + return delegate.tsMrevrange(range); + } + + @Override + public Flux> tsMrevrange(TimeRange range, MRangeOptions options) { + return delegate.tsMrevrange(range, options); + } + + @Override + public Mono tsGet(K key) { + return delegate.tsGet(key); + } + + @Override + public Flux> tsMget(MGetOptions options) { + return delegate.tsMget(options); + } + + @Override + public Flux> tsMget(V... filters) { + return delegate.tsMget(filters); + } + + @Override + public Flux> tsMgetWithLabels(V... filters) { + return delegate.tsMgetWithLabels(filters); + } + + @Override + public Flux tsInfo(K key) { + return delegate.tsInfo(key); + } + + @Override + public Flux tsInfoDebug(K key) { + return delegate.tsInfoDebug(key); + } + + @Override + public Flux tsQueryIndex(V... filters) { + return delegate.tsQueryIndex(filters); + } + + @Override + public Mono tsDel(K key, TimeRange timeRange) { + return delegate.tsDel(key, timeRange); + } + + @Override + public Mono bfAdd(K key, V item) { + return delegate.bfAdd(key, item); + } + + @Override + public Mono bfCard(K key) { + return delegate.bfCard(key); + } + + @Override + public Mono bfExists(K key, V item) { + return delegate.bfExists(key, item); + } + + @Override + public Mono bfInfo(K key) { + return delegate.bfInfo(key); + } + + @Override + public Mono bfInfo(K key, BloomFilterInfoType infoType) { + return delegate.bfInfo(key, infoType); + } + + @Override + public Flux bfInsert(K key, V... items) { + return delegate.bfInsert(key, items); + } + + @Override + public Flux bfInsert(K key, BloomFilterInsertOptions options, V... items) { + return delegate.bfInsert(key, options, items); + } + + @Override + public Flux bfMAdd(K key, V... items) { + return delegate.bfMAdd(key, items); + } + + @Override + public Flux bfMExists(K key, V... items) { + return delegate.bfMExists(key, items); + } + + @Override + public Mono bfReserve(K key, double errorRate, long capacity) { + return delegate.bfReserve(key, errorRate, capacity); + } + + @Override + public Mono bfReserve(K key, double errorRate, long capacity, BloomFilterReserveOptions options) { + return delegate.bfReserve(key, errorRate, capacity, options); + } + + @Override + public Mono cfAdd(K key, V item) { + return delegate.cfAdd(key, item); + } + + @Override + public Mono cfAddNx(K key, V item) { + return delegate.cfAddNx(key, item); + } + + @Override + public Mono cfCount(K key, V item) { + return delegate.cfCount(key, item); + } + + @Override + public Mono cfDel(K key, V item) { + return delegate.cfDel(key, item); + } + + @Override + public Mono cfExists(K key, V item) { + return delegate.cfExists(key, item); + } + + @Override + public Mono cfInfo(K key) { + return delegate.cfInfo(key); + } + + @Override + public Flux cfInsert(K key, V... items) { + return delegate.cfInsert(key, items); + } + + @Override + public Flux cfInsert(K key, CuckooFilterInsertOptions options, V... items) { + return delegate.cfInsert(key, options, items); + } + + @Override + public Flux cfInsertNx(K key, V... items) { + return delegate.cfInsertNx(key, items); + } + + @Override + public Flux cfInsertNx(K key, CuckooFilterInsertOptions options, V... items) { + return delegate.cfInsertNx(key, options, items); + } + + @Override + public Flux cfMExists(K key, V... items) { + return delegate.cfMExists(key, items); + } + + @Override + public Mono cfReserve(K key, long capacity) { + return delegate.cfReserve(key, capacity); + } + + @Override + public Mono cfReserve(K key, long capacity, CuckooFilterReserveOptions options) { + return delegate.cfReserve(key, capacity, options); + } + + @Override + public Mono cmsIncrBy(K key, V item, long increment) { + return delegate.cmsIncrBy(key, item, increment); + } + + @Override + public Flux cmsIncrBy(K key, LongScoredValue... itemIncrements) { + return delegate.cmsIncrBy(key, itemIncrements); + } + + @Override + public Mono cmsInitByProb(K key, double error, double probability) { + return delegate.cmsInitByProb(key, error, probability); + } + + @Override + public Mono cmsInitByDim(K key, long width, long depth) { + return delegate.cmsInitByDim(key, width, depth); + } + + @Override + public Flux cmsQuery(K key, V... items) { + return delegate.cmsQuery(key, items); + } + + @Override + public Mono cmsMerge(K destKey, K... keys) { + return delegate.cmsMerge(destKey, keys); + } + + @Override + public Mono cmsMerge(K destKey, LongScoredValue... sourceKeyWeights) { + return delegate.cmsMerge(destKey, sourceKeyWeights); + } + + @Override + public Mono cmsInfo(K key) { + return delegate.cmsInfo(key); + } + } diff --git a/core/lettucemod/src/main/java/com/redis/lettucemod/search/SearchCommandBuilder.java b/core/lettucemod/src/main/java/com/redis/lettucemod/search/SearchCommandBuilder.java index 3b095e1..fb1c0ff 100644 --- a/core/lettucemod/src/main/java/com/redis/lettucemod/search/SearchCommandBuilder.java +++ b/core/lettucemod/src/main/java/com/redis/lettucemod/search/SearchCommandBuilder.java @@ -1,6 +1,10 @@ package com.redis.lettucemod.search; +import java.util.Arrays; +import java.util.HashSet; import java.util.List; +import java.util.Set; +import java.util.function.Consumer; import com.redis.lettucemod.RedisModulesCommandBuilder; import com.redis.lettucemod.output.AggregateOutput; @@ -12,6 +16,7 @@ import com.redis.lettucemod.protocol.SearchCommandType; import io.lettuce.core.codec.RedisCodec; +import io.lettuce.core.codec.StringCodec; import io.lettuce.core.internal.LettuceAssert; import io.lettuce.core.output.BooleanOutput; import io.lettuce.core.output.CommandOutput; @@ -28,265 +33,295 @@ */ public class SearchCommandBuilder extends RedisModulesCommandBuilder { - public SearchCommandBuilder(RedisCodec codec) { - super(codec); - } - - protected Command createCommand(SearchCommandType type, CommandOutput output, - CommandArgs args) { - return new Command<>(type, output, args); - } - - private static void notNullIndex(Object index) { - notNull(index, "Index"); - } - - @SuppressWarnings("unchecked") - public Command create(K index, CreateOptions options, Field... fields) { - notNullIndex(index); - LettuceAssert.isTrue(fields.length > 0, "At least one field is required."); - SearchCommandArgs args = args(index); - if (options != null) { - options.build(args); - } - args.add(SearchCommandKeyword.SCHEMA); - for (Field field : fields) { - field.build(args); - } - return createCommand(SearchCommandType.CREATE, new StatusOutput<>(codec), args); - } - - public Command dropIndex(K index, boolean deleteDocs) { - notNullIndex(index); - SearchCommandArgs args = args(index); - if (deleteDocs) { - args.add(SearchCommandKeyword.DD); - } - return createCommand(SearchCommandType.DROPINDEX, new StatusOutput<>(codec), args); - } - - public Command> info(K index) { - notNullIndex(index); - SearchCommandArgs args = args(index); - return createCommand(SearchCommandType.INFO, new NestedMultiOutput<>(codec), args); - } - - public Command alter(K index, Field field) { - notNullIndex(index); - notNull(field, "Field"); - SearchCommandArgs args = args(index); - args.add(SearchCommandKeyword.SCHEMA); - args.add(SearchCommandKeyword.ADD); - field.build(args); - return createCommand(SearchCommandType.ALTER, new StatusOutput<>(codec), args); - } - - @Override - protected SearchCommandArgs args(K key) { - return new SearchCommandArgs<>(codec).addKey(key); - } - - private static void notNullQuery(Object query) { - notNull(query, "Query"); - } - - public Command> search(K index, V query, SearchOptions options) { - notNullIndex(index); - notNullQuery(query); - SearchCommandArgs args = args(index); - args.addValue(query); - if (options != null) { - options.build(args); - } - return createCommand(SearchCommandType.SEARCH, searchOutput(options), args); - } - - private CommandOutput> searchOutput(SearchOptions options) { - if (options == null) { - return new SearchOutput<>(codec); - } - if (options.isNoContent()) { - return new SearchNoContentOutput<>(codec, options.isWithScores()); - } - return new SearchOutput<>(codec, options.isWithScores(), options.isWithSortKeys(), options.isWithPayloads()); - } - - public Command> aggregate(K index, V query, AggregateOptions options) { - notNullIndex(index); - notNullQuery(query); - SearchCommandArgs args = args(index); - args.addValue(query); - if (options != null) { - options.build(args); - } - return createCommand(SearchCommandType.AGGREGATE, new AggregateOutput<>(codec, new AggregateResults<>()), args); - } - - public Command> aggregate(K index, V query, CursorOptions cursor, - AggregateOptions options) { - notNullIndex(index); - notNullQuery(query); - SearchCommandArgs args = args(index); - args.addValue(query); - if (options != null) { - options.build(args); - } - args.add(SearchCommandKeyword.WITHCURSOR); - if (cursor != null) { - cursor.build(args); - } - return createCommand(SearchCommandType.AGGREGATE, new AggregateWithCursorOutput<>(codec), args); - } - - public Command> cursorRead(K index, long cursor, Long count) { - notNullIndex(index); - SearchCommandArgs args = new SearchCommandArgs<>(codec); - args.add(SearchCommandKeyword.READ); - args.addKey(index); - args.add(cursor); - if (count != null) { - args.add(SearchCommandKeyword.COUNT); - args.add(count); - } - return createCommand(SearchCommandType.CURSOR, new AggregateWithCursorOutput<>(codec), args); - } - - public Command cursorDelete(K index, long cursor) { - notNullIndex(index); - SearchCommandArgs args = new SearchCommandArgs<>(codec); - args.add(SearchCommandKeyword.DEL); - args.addKey(index); - args.add(cursor); - return createCommand(SearchCommandType.CURSOR, new StatusOutput<>(codec), args); - } - - public Command> tagVals(K index, K field) { - notNullIndex(index); - SearchCommandArgs args = args(index); - args.addKey(field); - return createCommand(SearchCommandType.TAGVALS, new ValueListOutput<>(codec), args); - } - - private static void notNullDict(Object dict) { - notNull(dict, "Dict"); - } - - @SuppressWarnings("unchecked") - public Command dictadd(K dict, V... terms) { - notNullDict(dict); - return createCommand(SearchCommandType.DICTADD, new IntegerOutput<>(codec), args(dict).addValues(terms)); - } - - @SuppressWarnings("unchecked") - public Command dictdel(K dict, V... terms) { - notNullDict(dict); - return createCommand(SearchCommandType.DICTDEL, new IntegerOutput<>(codec), args(dict).addValues(terms)); - } - - public Command> dictdump(K dict) { - notNullDict(dict); - return createCommand(SearchCommandType.DICTDUMP, new ValueListOutput<>(codec), args(dict)); - } - - public Command sugadd(K key, V string, double score) { - return sugadd(key, string, score, null, false); - } - - public Command sugaddIncr(K key, V string, double score) { - return sugadd(key, string, score, null, true); - } - - public Command sugadd(K key, V string, double score, V payload) { - return sugadd(key, string, score, payload, false); - } - - public Command sugaddIncr(K key, V string, double score, V payload) { - return sugadd(key, string, score, payload, true); - } - - public Command sugadd(K key, V string, double score, V payload, boolean increment) { - notNullKey(key); - notNull(string, "String"); - SearchCommandArgs args = args(key); - args.addValue(string); - args.add(score); - if (increment) { - args.add(SearchCommandKeyword.INCR); - } - if (payload != null) { - args.add(SearchCommandKeyword.PAYLOAD); - args.addValue(payload); - } - return createCommand(SearchCommandType.SUGADD, new IntegerOutput<>(codec), args); - } - - public Command sugadd(K key, Suggestion suggestion) { - notNull(suggestion, "Suggestion"); - return sugadd(key, suggestion.getString(), suggestion.getScore(), suggestion.getPayload(), false); - } - - public Command sugaddIncr(K key, Suggestion suggestion) { - notNull(suggestion, "Suggestion"); - return sugadd(key, suggestion.getString(), suggestion.getScore(), suggestion.getPayload(), true); - } - - public Command>> sugget(K key, V prefix) { - return sugget(key, prefix, null); - } - - public Command>> sugget(K key, V prefix, SuggetOptions options) { - notNullKey(key); - notNull(prefix, "Prefix"); - SearchCommandArgs args = args(key); - args.addValue(prefix); - if (options != null) { - options.build(args); - } - return createCommand(SearchCommandType.SUGGET, suggetOutput(options), args); - } - - private SuggetOutput suggetOutput(SuggetOptions options) { - if (options == null) { - return new SuggetOutput<>(codec); - } - return new SuggetOutput<>(codec, options.isWithScores(), options.isWithPayloads()); - } - - public Command sugdel(K key, V string) { - notNullKey(key); - notNull(string, "String"); - return createCommand(SearchCommandType.SUGDEL, new BooleanOutput<>(codec), args(key).addValue(string)); - } - - public Command suglen(K key) { - notNullKey(key); - return createCommand(SearchCommandType.SUGLEN, new IntegerOutput<>(codec), args(key)); - - } - - private static void notNullName(Object name) { - notNull(name, "Name"); - } - - public Command aliasAdd(K name, K index) { - notNullName(name); - notNullIndex(index); - return createCommand(SearchCommandType.ALIASADD, new StatusOutput<>(codec), args(name).addKey(index)); - } - - public Command aliasUpdate(K name, K index) { - notNullName(name); - notNullIndex(index); - return createCommand(SearchCommandType.ALIASUPDATE, new StatusOutput<>(codec), args(name).addKey(index)); - } - - public Command aliasDel(K name) { - notNullName(name); - return createCommand(SearchCommandType.ALIASDEL, new StatusOutput<>(codec), args(name)); - } - - public Command> list() { - return new Command<>(SearchCommandType.LIST, new KeyListOutput<>(codec)); - } + public SearchCommandBuilder(RedisCodec codec) { + super(codec); + } + + protected Command createCommand(SearchCommandType type, CommandOutput output, + CommandArgs args) { + return new Command<>(type, output, args); + } + + private static void notNullIndex(Object index) { + notNull(index, "Index"); + } + + @SuppressWarnings("unchecked") + public Command create(K index, CreateOptions options, Field... fields) { + notNullIndex(index); + LettuceAssert.isTrue(fields.length > 0, "At least one field is required."); + SearchCommandArgs args = args(index); + if (options != null) { + options.build(args); + } + args.add(SearchCommandKeyword.SCHEMA); + for (Field field : fields) { + field.build(args); + } + return createCommand(SearchCommandType.CREATE, new StatusOutput<>(codec), args); + } + + public Command dropIndex(K index, boolean deleteDocs) { + notNullIndex(index); + SearchCommandArgs args = args(index); + if (deleteDocs) { + args.add(SearchCommandKeyword.DD); + } + return createCommand(SearchCommandType.DROPINDEX, new StatusOutput<>(codec), args); + } + + public Command> info(K index) { + notNullIndex(index); + SearchCommandArgs args = args(index); + return createCommand(SearchCommandType.INFO, new NestedMultiOutput<>(codec), args); + } + + public Command alter(K index, Field field) { + notNullIndex(index); + notNull(field, "Field"); + SearchCommandArgs args = args(index); + args.add(SearchCommandKeyword.SCHEMA); + args.add(SearchCommandKeyword.ADD); + field.build(args); + return createCommand(SearchCommandType.ALTER, new StatusOutput<>(codec), args); + } + + @Override + protected SearchCommandArgs args(K key) { + return new SearchCommandArgs<>(codec).addKey(key); + } + + private static void notNullQuery(Object query) { + notNull(query, "Query"); + } + + public Command> search(K index, V query, SearchOptions options) { + notNullIndex(index); + notNullQuery(query); + SearchCommandArgs args = args(index); + args.addValue(query); + if (options != null) { + options.build(args); + } + return createCommand(SearchCommandType.SEARCH, searchOutput(options), args); + } + + private CommandOutput> searchOutput(SearchOptions options) { + if (options == null) { + return new SearchOutput<>(codec); + } + if (options.isNoContent()) { + return new SearchNoContentOutput<>(codec, options.isWithScores()); + } + return new SearchOutput<>(codec, options.isWithScores(), options.isWithSortKeys(), options.isWithPayloads()); + } + + public Command> search(K index, V query, V... options) { + notNullIndex(index); + notNullQuery(query); + SearchCommandArgs args = args(index); + args.addValue(query); + SearchOptions searchOptions = new SearchOptions<>(); + args.addValues(options); + for (V option : options) { + String optionString = StringCodec.UTF8.decodeValue(codec.encodeValue(option)); + if (SearchCommandKeyword.WITHSORTKEYS.name().equalsIgnoreCase(optionString)) { + searchOptions.setWithSortKeys(true); + } + if (SearchCommandKeyword.WITHPAYLOADS.name().equalsIgnoreCase(optionString)) { + searchOptions.setWithPayloads(true); + } + if (SearchCommandKeyword.WITHSCORES.name().equalsIgnoreCase(optionString)) { + searchOptions.setWithScores(true); + } + } + return createCommand(SearchCommandType.SEARCH, searchOutput(searchOptions), args); + } + + public Command> aggregate(K index, V query, V... options) { + return doAggregate(index, query, args -> args.addValues(options)); + } + + public Command> aggregate(K index, V query, AggregateOptions options) { + return doAggregate(index, query, options == null ? null : options::build); + } + + private Command> doAggregate(K index, V query, Consumer> options) { + notNullIndex(index); + notNullQuery(query); + SearchCommandArgs args = args(index); + args.addValue(query); + if (options != null) { + options.accept(args); + } + return createCommand(SearchCommandType.AGGREGATE, new AggregateOutput<>(codec, new AggregateResults<>()), args); + } + + public Command> aggregate(K index, V query, CursorOptions cursor, + AggregateOptions options) { + notNullIndex(index); + notNullQuery(query); + SearchCommandArgs args = args(index); + args.addValue(query); + if (options != null) { + options.build(args); + } + args.add(SearchCommandKeyword.WITHCURSOR); + if (cursor != null) { + cursor.build(args); + } + return createCommand(SearchCommandType.AGGREGATE, new AggregateWithCursorOutput<>(codec), args); + } + + public Command> cursorRead(K index, long cursor, Long count) { + notNullIndex(index); + SearchCommandArgs args = new SearchCommandArgs<>(codec); + args.add(SearchCommandKeyword.READ); + args.addKey(index); + args.add(cursor); + if (count != null) { + args.add(SearchCommandKeyword.COUNT); + args.add(count); + } + return createCommand(SearchCommandType.CURSOR, new AggregateWithCursorOutput<>(codec), args); + } + + public Command cursorDelete(K index, long cursor) { + notNullIndex(index); + SearchCommandArgs args = new SearchCommandArgs<>(codec); + args.add(SearchCommandKeyword.DEL); + args.addKey(index); + args.add(cursor); + return createCommand(SearchCommandType.CURSOR, new StatusOutput<>(codec), args); + } + + public Command> tagVals(K index, K field) { + notNullIndex(index); + SearchCommandArgs args = args(index); + args.addKey(field); + return createCommand(SearchCommandType.TAGVALS, new ValueListOutput<>(codec), args); + } + + private static void notNullDict(Object dict) { + notNull(dict, "Dict"); + } + + @SuppressWarnings("unchecked") + public Command dictadd(K dict, V... terms) { + notNullDict(dict); + return createCommand(SearchCommandType.DICTADD, new IntegerOutput<>(codec), args(dict).addValues(terms)); + } + + @SuppressWarnings("unchecked") + public Command dictdel(K dict, V... terms) { + notNullDict(dict); + return createCommand(SearchCommandType.DICTDEL, new IntegerOutput<>(codec), args(dict).addValues(terms)); + } + + public Command> dictdump(K dict) { + notNullDict(dict); + return createCommand(SearchCommandType.DICTDUMP, new ValueListOutput<>(codec), args(dict)); + } + + public Command sugadd(K key, V string, double score) { + return sugadd(key, string, score, null, false); + } + + public Command sugaddIncr(K key, V string, double score) { + return sugadd(key, string, score, null, true); + } + + public Command sugadd(K key, V string, double score, V payload) { + return sugadd(key, string, score, payload, false); + } + + public Command sugaddIncr(K key, V string, double score, V payload) { + return sugadd(key, string, score, payload, true); + } + + public Command sugadd(K key, V string, double score, V payload, boolean increment) { + notNullKey(key); + notNull(string, "String"); + SearchCommandArgs args = args(key); + args.addValue(string); + args.add(score); + if (increment) { + args.add(SearchCommandKeyword.INCR); + } + if (payload != null) { + args.add(SearchCommandKeyword.PAYLOAD); + args.addValue(payload); + } + return createCommand(SearchCommandType.SUGADD, new IntegerOutput<>(codec), args); + } + + public Command sugadd(K key, Suggestion suggestion) { + notNull(suggestion, "Suggestion"); + return sugadd(key, suggestion.getString(), suggestion.getScore(), suggestion.getPayload(), false); + } + + public Command sugaddIncr(K key, Suggestion suggestion) { + notNull(suggestion, "Suggestion"); + return sugadd(key, suggestion.getString(), suggestion.getScore(), suggestion.getPayload(), true); + } + + public Command>> sugget(K key, V prefix) { + return sugget(key, prefix, null); + } + + public Command>> sugget(K key, V prefix, SuggetOptions options) { + notNullKey(key); + notNull(prefix, "Prefix"); + SearchCommandArgs args = args(key); + args.addValue(prefix); + if (options != null) { + options.build(args); + } + return createCommand(SearchCommandType.SUGGET, suggetOutput(options), args); + } + + private SuggetOutput suggetOutput(SuggetOptions options) { + if (options == null) { + return new SuggetOutput<>(codec); + } + return new SuggetOutput<>(codec, options.isWithScores(), options.isWithPayloads()); + } + + public Command sugdel(K key, V string) { + notNullKey(key); + notNull(string, "String"); + return createCommand(SearchCommandType.SUGDEL, new BooleanOutput<>(codec), args(key).addValue(string)); + } + + public Command suglen(K key) { + notNullKey(key); + return createCommand(SearchCommandType.SUGLEN, new IntegerOutput<>(codec), args(key)); + + } + + private static void notNullName(Object name) { + notNull(name, "Name"); + } + + public Command aliasAdd(K name, K index) { + notNullName(name); + notNullIndex(index); + return createCommand(SearchCommandType.ALIASADD, new StatusOutput<>(codec), args(name).addKey(index)); + } + + public Command aliasUpdate(K name, K index) { + notNullName(name); + notNullIndex(index); + return createCommand(SearchCommandType.ALIASUPDATE, new StatusOutput<>(codec), args(name).addKey(index)); + } + + public Command aliasDel(K name) { + notNullName(name); + return createCommand(SearchCommandType.ALIASDEL, new StatusOutput<>(codec), args(name)); + } + + public Command> list() { + return new Command<>(SearchCommandType.LIST, new KeyListOutput<>(codec)); + } } diff --git a/core/lettucemod/src/test/java/com/redis/lettucemod/ModulesTests.java b/core/lettucemod/src/test/java/com/redis/lettucemod/ModulesTests.java index 5fc624a..c1dc79c 100644 --- a/core/lettucemod/src/test/java/com/redis/lettucemod/ModulesTests.java +++ b/core/lettucemod/src/test/java/com/redis/lettucemod/ModulesTests.java @@ -334,8 +334,20 @@ void ftSearchOptions() throws IOException { Document doc1 = results.get(0); assertNotNull(doc1.get(NAME)); assertNotNull(doc1.get(STYLE)); + assertEquals(doc1.get(DESCRIPTION), doc1.getPayload()); assertNull(abv(doc1)); + } + @Test + void ftSearchStringOptions() throws IOException { + populateIndex(); + String options = "withpayloads limit 0 10 WITHSCORES nostopwords"; + SearchResults results = commands.ftSearch(INDEX, "pale", options.split(" ")); + assertEquals(74, results.getCount()); + Document doc1 = results.get(0); + assertNotNull(doc1.get(NAME)); + assertNotNull(doc1.get(STYLE)); + assertEquals(doc1.get(DESCRIPTION), doc1.getPayload()); } private void assertSearch(String query, SearchOptions options, long expectedCount, String... expectedAbv) { @@ -389,8 +401,9 @@ void ftSearchTags() throws InterruptedException, ExecutionException, IOException String term = "pale"; String query = "@style:" + term; Tags tags = new Tags<>("", ""); - SearchResults results = commands.ftSearch(INDEX, query, SearchOptions. builder() - .highlight(SearchOptions.Highlight. builder().build()).build()); + SearchResults results = commands.ftSearch(INDEX, query, + SearchOptions. builder().highlight(SearchOptions.Highlight. builder().build()) + .build()); for (Document result : results) { assertTrue(isHighlighted(result, STYLE, tags, term)); } @@ -493,13 +506,27 @@ void ftAggregateGroupSortLimit() throws Exception { assertTrue(abvs.get(0) > abvs.get(abvs.size() - 1)); assertEquals(20, results.size()); }; - AggregateOptions options = AggregateOptions - . operation(Group.by(STYLE).reducer(Avg.property(ABV).as(ABV).build()).build()) + AggregateOptions options = AggregateOptions. operation( + Group.by(STYLE).reducer(Avg.property(ABV).as(ABV).build()).build()) .operation(Sort.by(Sort.Property.desc(ABV)).build()).operation(Limit.offset(0).num(20)).build(); asserts.accept(commands.ftAggregate(INDEX, "*", options)); asserts.accept(connection.reactive().ftAggregate(INDEX, "*", options).block()); } + @Test + void ftAggregateStringOptions() throws Exception { + populateBeers(); + Consumer> asserts = results -> { + assertEquals(36, results.getCount()); + List abvs = results.stream().map(r -> Double.parseDouble((String) r.get(ABV))).collect(Collectors.toList()); + assertTrue(abvs.get(0) > abvs.get(abvs.size() - 1)); + assertEquals(20, results.size()); + }; + String[] options = "groupby 1 @style_name reduce avg 1 @abv as abv sortby 2 @abv desc limit 0 20".split(" "); + asserts.accept(commands.ftAggregate(INDEX, "*", options)); + asserts.accept(connection.reactive().ftAggregate(INDEX, "*", options).block()); + } + @Test void ftAggregateSort() throws Exception { populateBeers(); @@ -509,8 +536,8 @@ void ftAggregateSort() throws Exception { List abvs = results.stream().map(this::abv).collect(Collectors.toList()); assertTrue(abvs.get(0) > abvs.get(abvs.size() - 1)); }; - AggregateOptions options = AggregateOptions - . operation(Sort.by(Sort.Property.desc(ABV)).by(Property.asc(IBU)).build()).build(); + AggregateOptions options = AggregateOptions. operation( + Sort.by(Sort.Property.desc(ABV)).by(Property.asc(IBU)).build()).build(); asserts.accept(commands.ftAggregate(INDEX, "*", options)); asserts.accept(connection.reactive().ftAggregate(INDEX, "*", options).block()); } @@ -533,8 +560,8 @@ void ftAggregateGroupNone() throws Exception { assertEquals(16, maxAbv, 0.1); }; - AggregateOptions options = AggregateOptions - . operation(Group.by().reducer(Max.property(ABV).as(ABV).build()).build()).build(); + AggregateOptions options = AggregateOptions. operation( + Group.by().reducer(Max.property(ABV).as(ABV).build()).build()).build(); asserts.accept(commands.ftAggregate(INDEX, "*", options)); asserts.accept(connection.reactive().ftAggregate(INDEX, "*", options).block()); @@ -574,8 +601,9 @@ void ftAggregateCursor() throws Exception { AggregateWithCursorResults cursorResults = commands.ftAggregate(INDEX, "*", CursorOptions.builder().count(10).build(), cursorOptions); cursorTests.accept(cursorResults); - cursorTests.accept(connection.reactive() - .ftAggregate(INDEX, "*", CursorOptions.builder().count(10).build(), cursorOptions).block()); + cursorTests.accept( + connection.reactive().ftAggregate(INDEX, "*", CursorOptions.builder().count(10).build(), cursorOptions) + .block()); cursorResults = commands.ftCursorRead(INDEX, cursorResults.getCursor(), 400); assertEquals(400, cursorResults.size()); String deleteStatus = commands.ftCursorDelete(INDEX, cursorResults.getCursor()); @@ -690,8 +718,8 @@ private String jsonField(String name) { void ftTagVals() throws Exception { populateIndex(); Set TAG_VALS = new HashSet<>(Arrays.asList( - "american-style brown ale, traditional german-style bock, german-style schwarzbier, old ale, american-style india pale ale, german-style oktoberfest, other belgian-style ales, american-style stout, winter warmer, belgian-style tripel, american-style lager, belgian-style dubbel, porter, american-style barley wine ale, belgian-style fruit lambic, scottish-style light ale, south german-style hefeweizen, imperial or double india pale ale, golden or blonde ale, belgian-style quadrupel, american-style imperial stout, belgian-style pale strong ale, english-style pale mild ale, american-style pale ale, irish-style red ale, dark american-belgo-style ale, light american wheat ale or lager, german-style pilsener, american-style amber/red ale, scotch ale, german-style doppelbock, extra special bitter, south german-style weizenbock, english-style india pale ale, belgian-style pale ale, french & belgian-style saison" - .split(", "))); + "american-style brown ale, traditional german-style bock, german-style schwarzbier, old ale, american-style india pale ale, german-style oktoberfest, other belgian-style ales, american-style stout, winter warmer, belgian-style tripel, american-style lager, belgian-style dubbel, porter, american-style barley wine ale, belgian-style fruit lambic, scottish-style light ale, south german-style hefeweizen, imperial or double india pale ale, golden or blonde ale, belgian-style quadrupel, american-style imperial stout, belgian-style pale strong ale, english-style pale mild ale, american-style pale ale, irish-style red ale, dark american-belgo-style ale, light american wheat ale or lager, german-style pilsener, american-style amber/red ale, scotch ale, german-style doppelbock, extra special bitter, south german-style weizenbock, english-style india pale ale, belgian-style pale ale, french & belgian-style saison".split( + ", "))); HashSet actual = new HashSet<>(commands.ftTagvals(INDEX, STYLE)); assertEquals(TAG_VALS, actual); assertEquals(TAG_VALS, new HashSet<>(connection.reactive().ftTagvals(INDEX, STYLE).collectList().block())); @@ -709,7 +737,7 @@ void ftAggregateEmptyToListReducer() { Map doc1 = mapOf("category", "31", "color", "red"); commands.hset("my_prefix:1", doc1); AggregateOptions aggregateOptions = AggregateOptions. operation(Group.by("category") - .reducers(ToList.property("color").as("color").build(), ToList.property("size").as("size").build()).build()) + .reducers(ToList.property("color").as("color").build(), ToList.property("size").as("size").build()).build()) .build(); AggregateResults results = commands.ftAggregate("idx", "@color:{red|blue}", aggregateOptions); assertEquals(1, results.size());