diff --git a/builtin-adapter/BuiltinAdapter.cpp b/builtin-adapter/BuiltinAdapter.cpp index 4bcd914d..eac37126 100644 --- a/builtin-adapter/BuiltinAdapter.cpp +++ b/builtin-adapter/BuiltinAdapter.cpp @@ -102,7 +102,9 @@ JNIEXPORT jobject JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinA reindexer_string dsn = rx_string(env, path); reindexer_string vers = rx_string(env, version); reindexer_error error = reindexer_connect(rx, dsn, ConnectOpts(), vers, BindingCapabilities( - kBindingCapabilityResultsWithShardIDs | kBindingCapabilityComplexRank)); + kBindingCapabilityResultsWithShardIDs + | kBindingCapabilityComplexRank + | kBindingCapabilityQueryFormatV2)); env->ReleaseStringUTFChars(path, reinterpret_cast(dsn.p)); env->ReleaseStringUTFChars(version, reinterpret_cast(vers.p)); return j_res(env, error); diff --git a/src/main/java/ru/rt/restream/reindexer/Query.java b/src/main/java/ru/rt/restream/reindexer/Query.java index 45fdabee..e2885835 100644 --- a/src/main/java/ru/rt/restream/reindexer/Query.java +++ b/src/main/java/ru/rt/restream/reindexer/Query.java @@ -34,6 +34,7 @@ import java.util.ArrayDeque; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.Deque; @@ -58,6 +59,8 @@ import static ru.rt.restream.reindexer.binding.Consts.MERGE; import static ru.rt.restream.reindexer.binding.Consts.MODE_ACCURATE_TOTAL; import static ru.rt.restream.reindexer.binding.Consts.OR_INNER_JOIN; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V1; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2; import static ru.rt.restream.reindexer.binding.Consts.VALUE_STRING; /** @@ -210,11 +213,15 @@ public enum Condition { private Query root; + private final int queryFormatVersion; + Query(Reindexer reindexer, ReindexerNamespace namespace, TransactionContext transactionContext) { logBuilder.namespace(namespace.getName()); this.reindexer = reindexer; this.namespace = namespace; this.transactionContext = transactionContext; + this.queryFormatVersion = reindexer.getBinding().queryFormatVersion(); + buffer.putUInt8(0); buffer.putVString(namespace.getName()); } @@ -437,7 +444,7 @@ public Query where(Query subquery, Condition condition, Object... values) logBuilder.where(nextOperation, subquery, condition.code, values); buffer.putVarUInt32(QUERY_SUB_QUERY_CONDITION) .putVarUInt32(nextOperation) - .putVBytes(subquery.buffer.bytes()) + .putVBytes(subquery.bytes()) .putVarUInt32(condition.code); this.nextOperation = OP_AND; @@ -465,7 +472,7 @@ public Query where(String indexName, Condition condition, Query subquery) .putVarUInt32(nextOperation) .putVString(indexName) .putVarUInt32(condition.code) - .putVBytes(subquery.buffer.bytes()); + .putVBytes(subquery.bytes()); this.nextOperation = OP_AND; this.queryCount++; @@ -967,11 +974,12 @@ public ResultIterator execute() { * @return an iterator over a query result */ public ResultIterator execute(Class itemClass) { - long[] ptVersions = prepareQueryAndGetPayloadTypesVersions(); + byte[] queryData = buildSelectQueryBytes(); + long[] payloadTypeVersions = getPayloadTypeVersions(); RequestContext requestContext = transactionContext != null - ? transactionContext.selectQuery(buffer.bytes(), fetchCount, ptVersions, false) - : reindexer.getBinding().selectQuery(buffer.bytes(), fetchCount, ptVersions, false); + ? transactionContext.selectQuery(queryData, fetchCount, payloadTypeVersions, false) + : reindexer.getBinding().selectQuery(queryData, fetchCount, payloadTypeVersions, false); updatePayloadTypes(requestContext.getQueryResult()); @@ -984,11 +992,12 @@ public ResultIterator execute(Class itemClass) { * @return an iterator over a query result */ public QueryResultJsonIterator executeToJson() { - long[] ptVersions = prepareQueryAndGetPayloadTypesVersions(); + byte[] queryData = buildSelectQueryBytes(); + long[] payloadTypeVersions = getPayloadTypeVersions(); RequestContext requestContext = transactionContext != null - ? transactionContext.selectQuery(buffer.bytes(), fetchCount, ptVersions, true) - : reindexer.getBinding().selectQuery(buffer.bytes(), fetchCount, ptVersions, true); + ? transactionContext.selectQuery(queryData, fetchCount, payloadTypeVersions, true) + : reindexer.getBinding().selectQuery(queryData, fetchCount, payloadTypeVersions, true); QueryResult queryResult = requestContext.getQueryResult(); @@ -1026,49 +1035,35 @@ private void updatePayloadTypes(QueryResult queryResult) { } } - private long[] prepareQueryAndGetPayloadTypesVersions() { + private byte[] buildSelectQueryBytes() { logBuilder.type(SELECT); if (LOGGER.isDebugEnabled()) { debug(); LOGGER.debug(logBuilder.getSql()); } + namespaces.clear(); namespaces.add(namespace); for (Query mergeQuery : mergeQueries) { namespaces.add(mergeQuery.namespace); } - for (Query joinQuery : joinQueries) { - namespaces.add(joinQuery.namespace); - } - - for (Query mergeQuery : mergeQueries) { - for (Query joinQuery : mergeQuery.joinQueries) { - namespaces.add(joinQuery.namespace); - } - } - - buffer.putVarUInt32(QUERY_END); - - for (Query joinQuery : joinQueries) { - buffer.putVarUInt32(joinQuery.joinType); - buffer.writeBytes(joinQuery.buffer.bytes()); - buffer.putVarUInt32(QUERY_END); + int formatVersion = reindexer.getBinding().queryFormatVersion(); + ByteBuffer queryBuffer = new ByteBuffer(getQueryBytes(formatVersion)); + queryBuffer.putVarUInt32(QUERY_END); + if (formatVersion == QUERY_FORMAT_V2) { + appendJoinQueries(queryBuffer, namespaces, formatVersion); + appendMergeQueries(queryBuffer, namespaces, formatVersion); + } else { + appendJoinQueriesV1(queryBuffer, true); + appendMergeQueriesV1(queryBuffer); } - for (Query mergeQuery : mergeQueries) { - buffer.putVarUInt32(MERGE); - buffer.writeBytes(mergeQuery.buffer.bytes()); - buffer.putVarUInt32(QUERY_END); - List> joinQueries = mergeQuery.getJoinQueries(); - for (Query joinQuery : joinQueries) { - buffer.putVarUInt32(joinQuery.joinType); - buffer.writeBytes(joinQuery.buffer.bytes()); - buffer.putVarUInt32(QUERY_END); - } - } + return queryBuffer.bytes(); + } + private long[] getPayloadTypeVersions() { return namespaces.stream() .map(ReindexerNamespace::getPayloadType) .mapToLong(pt -> pt == null ? 0 : (pt.getVersion() ^ pt.getStateToken())) @@ -1085,9 +1080,9 @@ public void delete() { LOGGER.debug(logBuilder.getSql()); } if (transactionContext != null) { - transactionContext.deleteQuery(buffer.bytes()); + transactionContext.deleteQuery(toExecutableBytes()); } else { - reindexer.getBinding().deleteQuery(buffer.bytes()); + reindexer.getBinding().deleteQuery(toExecutableBytes()); } } @@ -1263,13 +1258,13 @@ public void update() { LOGGER.debug(logBuilder.getSql()); } if (transactionContext != null) { - transactionContext.updateQuery(buffer.bytes()); + transactionContext.updateQuery(toExecutableBytes()); } else { // There are no support for inner joins for update-queries in Java binding, // so we are using single pt version PayloadType pt = namespace.getPayloadType(); long tmVersion = pt == null ? 0 : (pt.getVersion() ^ pt.getStateToken()); - reindexer.getBinding().updateQuery(buffer.bytes(), new long[]{tmVersion}); + reindexer.getBinding().updateQuery(toExecutableBytes(), new long[]{tmVersion}); } } @@ -1287,6 +1282,13 @@ public List> getMergeQueries() { return mergeQueries; } + /** + * Return query namespace. + */ + ReindexerNamespace getNamespace() { + return namespace; + } + /** * Get the namespace names that are used in the current query. */ @@ -1310,13 +1312,116 @@ String getSql() { return logBuilder.getSql(); } - /** - * Returns all used bytes from the {@link ByteBuffer}. - * - * @return all used bytes from the {@code ByteBuffer} - */ public byte[] bytes() { - return buffer.bytes(); + return toSubQueryBytes(queryFormatVersion); + } + + private byte[] toSubQueryBytes(int formatVersion) { + byte[] queryBytes = getQueryBytes(formatVersion); + if (formatVersion == QUERY_FORMAT_V2 || hasNestedJoins()) { + ByteBuffer copy = new ByteBuffer(queryBytes); + copy.putVarUInt32(QUERY_END); + copy.putVarUInt32(0); + copy.putVarUInt32(0); + return copy.bytes(); + } + return queryBytes; + } + + private byte[] toExecutableBytes() { + int formatVersion = reindexer.getBinding().queryFormatVersion(); + ByteBuffer queryBuffer = new ByteBuffer(getQueryBytes(formatVersion)); + queryBuffer.putVarUInt32(QUERY_END); + if (formatVersion == QUERY_FORMAT_V2) { + appendJoinQueries(queryBuffer, new ArrayList<>(), formatVersion); + appendMergeQueries(queryBuffer, new ArrayList<>(), formatVersion); + } else { + appendJoinQueriesV1(queryBuffer, false); + } + return queryBuffer.bytes(); + } + + private byte[] getQueryBytes(int formatVersion) { + byte[] queryBytes = buffer.bytes(); + if (queryBytes.length == 0) { + return queryBytes; + } + if (formatVersion == QUERY_FORMAT_V2) { + queryBytes[0] = (byte) QUERY_FORMAT_V2; + return queryBytes; + } + return Arrays.copyOfRange(queryBytes, 1, queryBytes.length); + } + + private void appendJoinQueries(ByteBuffer target, List> targetNamespaces, int formatVersion) { + target.putVarUInt32(joinQueries.size()); + for (Query joinQuery : joinQueries) { + appendQuery(target, joinQuery, joinQuery.joinType, targetNamespaces, formatVersion); + } + } + + private void appendMergeQueries(ByteBuffer target, List> targetNamespaces, int formatVersion) { + target.putVarUInt32(mergeQueries.size()); + for (Query mergeQuery : mergeQueries) { + appendQuery(target, mergeQuery, MERGE, targetNamespaces, formatVersion); + } + } + + private void appendQuery(ByteBuffer target, Query query, int queryJoinType, + List> targetNamespaces, int formatVersion) { + if (queryJoinType != MERGE) { + targetNamespaces.add(query.namespace); + } + + target.putVarUInt32(queryJoinType); + target.writeBytes(query.getQueryBytes(formatVersion)); + target.putVarUInt32(QUERY_END); + + query.appendJoinQueries(target, targetNamespaces, formatVersion); + query.appendMergeQueries(target, targetNamespaces, formatVersion); + } + + private void appendJoinQueriesV1(ByteBuffer target, boolean appendNamespaces) { + if (hasNestedJoins()) { + throw new IllegalStateException("Nested joins are not supported by QueryFormatV1"); + } + for (Query joinQuery : joinQueries) { + if (appendNamespaces) { + namespaces.add(joinQuery.namespace); + } + target.putVarUInt32(joinQuery.joinType); + target.writeBytes(joinQuery.getQueryBytes(QUERY_FORMAT_V1)); + target.putVarUInt32(QUERY_END); + } + } + + private void appendMergeQueriesV1(ByteBuffer target) { + for (Query mergeQuery : mergeQueries) { + target.putVarUInt32(MERGE); + target.writeBytes(mergeQuery.getQueryBytes(QUERY_FORMAT_V1)); + target.putVarUInt32(QUERY_END); + + for (Query joinQuery : mergeQuery.joinQueries) { + namespaces.add(joinQuery.namespace); + target.putVarUInt32(joinQuery.joinType); + target.writeBytes(joinQuery.getQueryBytes(QUERY_FORMAT_V1)); + target.putVarUInt32(QUERY_END); + } + } + } + + private boolean hasNestedJoins() { + for (Query joinQuery : joinQueries) { + if (!joinQuery.joinQueries.isEmpty() || joinQuery.hasNestedJoins()) { + return true; + } + } + for (Query mergeQuery : mergeQueries) { + if (mergeQuery.hasNestedJoins()) { + return true; + } + } + return false; } /** diff --git a/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java b/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java index 41eb1009..b51d54f9 100644 --- a/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java +++ b/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java @@ -123,97 +123,123 @@ public T next() { fetchResults(); } + T item = itemClass.cast(readItem(namespace, itemReader, query)); + position++; + return item; + + } + + private S readItem(ReindexerNamespace expectedNamespace, ItemReader reader, Query queryContext) { ItemParams params = readItemParams(); - T item; + Query itemQueryContext = getItemQueryContext(queryContext, params.nsId); + + ReindexerNamespace itemNamespace = expectedNamespace; + if (query != null && params.nsId < query.getNamespaces().size()) { + itemNamespace = query.getNamespaces().get(params.nsId); + } + + S item = readItemData(params, reader, itemNamespace); + readJoinedItems(item, itemQueryContext, params.nsId); + return item; + } + + private Query getItemQueryContext(Query defaultQueryContext, int nsId) { + if (query == null || nsId == 0) { + return defaultQueryContext; + } + + List> mergeQueries = query.getMergeQueries(); + if (nsId <= mergeQueries.size()) { + return mergeQueries.get(nsId - 1); + } + + return defaultQueryContext; + } + + private S readItemData(ItemParams params, ItemReader reader, ReindexerNamespace itemNamespace) { if (params.cptr != 0) { - ByteBuffer buffer = NativeUtils.getNativeBuffer(queryResult.getResultsPtr(), params.cptr, params.nsId); - item = itemReader.readItem(buffer); - } else { - int length = (int) buffer.getUInt32(); - item = itemReader.readItem(new ByteBuffer(buffer.getBytes(length)).rewind()); + ByteBuffer nativeBuffer = NativeUtils.getNativeBuffer(queryResult.getResultsPtr(), params.cptr, + params.nsId); + return reader.readItem(nativeBuffer); } - long subNsRes = -1L; - if (queryResult.isWithJoined()) { - subNsRes = buffer.getVarUInt(); + int length = (int) buffer.getUInt32(); + return reader.readItem(new ByteBuffer(buffer.getBytes(length)).rewind()); + } + + private void readJoinedItems(Object item, Query queryContext, int nsId) { + if (!queryResult.isWithJoined()) { + return; } - int nsIndexOffset = getJoinedNsIndexOffset(params.nsId); + if (queryResult.getQueryFormatVersion() == Consts.QUERY_FORMAT_V1) { + readJoinedItemsV1(item, nsId); + return; + } + int joinedFields = (int) buffer.getVarUInt(); Map> subItemsMap = new HashMap<>(); - for (int nsIndex = 0; nsIndex < subNsRes; nsIndex++) { - if (query == null) { - skipSubItems(); - } else { - readSubItems(nsIndexOffset, subItemsMap, nsIndex); + for (int joinedField = 0; joinedField < joinedFields; joinedField++) { + int itemsCount = (int) buffer.getVarUInt(); + if (queryContext == null) { + skipJoinedItems(itemsCount); + continue; } - } + Query joinQuery = queryContext.getJoinQueries().get(joinedField); + ReindexerNamespace joinedNamespace = joinQuery.getNamespace(); + CjsonItemReader joinedItemReader = newItemReader(joinedNamespace); + List subItems = new ArrayList<>(itemsCount); + for (int i = 0; i < itemsCount; i++) { + subItems.add(readItem(joinedNamespace, joinedItemReader, joinQuery)); + } + subItemsMap.computeIfAbsent(queryContext.getJoinFields().get(joinedField), field -> new ArrayList<>()) + .addAll(subItems); + } subItemsMap.forEach((key, value) -> writeJoinResult(item, key, value)); - - position++; - return item; - } - private void readSubItems(int nsIndexOffset, Map> subItemsMap, int nsIndex) { - List> namespaces = query.getNamespaces(); - int nsId = nsIndex + nsIndexOffset; - ReindexerNamespace subItemNamespace = namespaces.get(nsId); - PayloadType subItemPayloadType = subItemNamespace.getPayloadType(); - CtagMatcher ctagMatcher = new CtagMatcher(); - ctagMatcher.read(subItemPayloadType); - Class siClass = subItemNamespace.getItemClass(); - CjsonItemReader subItemItemReader = new CjsonItemReader<>(siClass, ctagMatcher); - String joinField = query.getJoinFields().get(nsIndex); - List subItems = subItemsMap.computeIfAbsent(joinField, s -> new ArrayList<>()); - - int siRes = (int) buffer.getVarUInt(); - for (int i = 0; i < siRes; i++) { - ItemParams subItemParams = readItemParams(); - Object subItem; - if (subItemParams.cptr != 0) { - ByteBuffer buffer = NativeUtils.getNativeBuffer(queryResult.getResultsPtr(), subItemParams.cptr, - nsId); - subItem = subItemItemReader.readItem(buffer); - } else { - int subItemLength = (int) buffer.getUInt32(); - subItem = subItemItemReader.readItem(new ByteBuffer(buffer.getBytes(subItemLength)).rewind()); + private void readJoinedItemsV1(Object item, int nsId) { + int joinedFields = (int) buffer.getVarUInt(); + if (query == null) { + for (int joinedField = 0; joinedField < joinedFields; joinedField++) { + skipJoinedItemsV1((int) buffer.getVarUInt()); } - subItems.add(subItem); + return; } - } - private void skipSubItems() { - int siRes = (int) buffer.getVarUInt(); - for (int i = 0; i < siRes; i++) { - ItemParams subItemParams = readItemParams(); - if (subItemParams.cptr == 0) { - int subItemLength = (int) buffer.getUInt32(); - buffer.skip(subItemLength); + int namespaceIndexOffset = getJoinedNsIndexOffset(nsId); + for (int nsIndex = 0; nsIndex < joinedFields; nsIndex++) { + int itemsCount = (int) buffer.getVarUInt(); + ReindexerNamespace joinedNamespace = query.getNamespaces().get(nsIndex + namespaceIndexOffset); + CjsonItemReader joinedItemReader = newItemReader(joinedNamespace); + List subItems = new ArrayList<>(itemsCount); + for (int j = 0; j < itemsCount; j++) { + ItemParams subItemParams = readItemParams(); + subItems.add(readItemData(subItemParams, joinedItemReader, joinedNamespace)); } + + writeJoinResult(item, query.getJoinFields().get(nsIndex), subItems); } } - private void writeJoinResult(T item, String fieldName, List subItems) { - Field field = FieldUtils.getField(item.getClass(), fieldName, true); - - if (field == null || !field.isAnnotationPresent(Transient.class)) { - String msg = String.format("Join results omitted: no transient field '%s' found", fieldName); - LOGGER.debug(msg); + private void skipJoinedItems(int itemsCount) { + for (int i = 0; i < itemsCount; i++) { + ItemParams itemParams = readItemParams(); + if (itemParams.cptr == 0) { + int length = (int) buffer.getUInt32(); + buffer.skip(length); + } + readJoinedItems(null, null, itemParams.nsId); } + } - if (field != null) { - if (field.getType() == List.class) { - BeanPropertyUtils.setProperty(item, fieldName, subItems); - } else { - if (subItems.size() > 1) { - throw new RuntimeException("Multiple join result found: " + fieldName); - } else if (subItems.size() == 0) { - BeanPropertyUtils.setProperty(item, fieldName, null); - } else { - BeanPropertyUtils.setProperty(item, fieldName, subItems.get(0)); - } + private void skipJoinedItemsV1(int itemsCount) { + for (int i = 0; i < itemsCount; i++) { + ItemParams itemParams = readItemParams(); + if (itemParams.cptr == 0) { + int length = (int) buffer.getUInt32(); + buffer.skip(length); } } } @@ -238,6 +264,36 @@ private int getJoinedNsIndexOffset(int nsId) { return offset; } + private CjsonItemReader newItemReader(ReindexerNamespace itemNamespace) { + PayloadType payloadType = itemNamespace.getPayloadType(); + CtagMatcher ctagMatcher = new CtagMatcher(); + ctagMatcher.read(payloadType); + return new CjsonItemReader<>(itemNamespace.getItemClass(), ctagMatcher); + } + + private void writeJoinResult(Object item, String fieldName, List subItems) { + Field field = FieldUtils.getField(item.getClass(), fieldName, true); + + if (field == null || !field.isAnnotationPresent(Transient.class)) { + String msg = String.format("Join results omitted: no transient field '%s' found", fieldName); + LOGGER.debug(msg); + } + + if (field != null) { + if (field.getType() == List.class) { + BeanPropertyUtils.setProperty(item, fieldName, subItems); + } else { + if (subItems.size() > 1) { + throw new RuntimeException("Multiple join result found: " + fieldName); + } else if (subItems.size() == 0) { + BeanPropertyUtils.setProperty(item, fieldName, null); + } else { + BeanPropertyUtils.setProperty(item, fieldName, subItems.get(0)); + } + } + } + } + private ItemParams readItemParams() { ItemParams params = new ItemParams(); diff --git a/src/main/java/ru/rt/restream/reindexer/binding/Binding.java b/src/main/java/ru/rt/restream/reindexer/binding/Binding.java index e33f7a5e..50c40299 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/Binding.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/Binding.java @@ -205,6 +205,15 @@ public interface Binding { */ String getMeta(String namespaceName, String key); + /** + * Returns negotiated query serialization format version. + * + * @return query serialization format version + */ + default int queryFormatVersion() { + return Consts.QUERY_FORMAT_V1; + } + /** * Closes binding to Reindexer instance. */ diff --git a/src/main/java/ru/rt/restream/reindexer/binding/Consts.java b/src/main/java/ru/rt/restream/reindexer/binding/Consts.java index 56044f6e..63fec6bc 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/Consts.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/Consts.java @@ -120,6 +120,10 @@ public final class Consts { // incarnation tags are not supported in java connector public static final long BINDING_CAPABILITY_NAMESPACE_INCARNATIONS = 1 << 2; public static final long BINDING_CAPABILITY_COMPLEX_RANK = 1 << 3; + public static final long BINDING_CAPABILITY_QUERY_FORMAT_V2 = 1 << 4; + + public static final int QUERY_FORMAT_V1 = 1; + public static final int QUERY_FORMAT_V2 = 2; public static final int RANK_FORMAT_SINGLE_FLOAT = 0; public static final float EMPTY_RANK = Float.NEGATIVE_INFINITY; diff --git a/src/main/java/ru/rt/restream/reindexer/binding/QueryResult.java b/src/main/java/ru/rt/restream/reindexer/binding/QueryResult.java index 5ff24fe9..5d0a2c17 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/QueryResult.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/QueryResult.java @@ -122,4 +122,9 @@ public class QueryResult { */ private long rankFormat; + /** + * Query serialization format version used for this result. + */ + private int queryFormatVersion = Consts.QUERY_FORMAT_V1; + } diff --git a/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java b/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java index 0b32b0fd..62aa2779 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java @@ -33,6 +33,7 @@ import static ru.rt.restream.reindexer.binding.Consts.QUERY_RESULT_RANK_FORMAT; import static ru.rt.restream.reindexer.binding.Consts.QUERY_RESULT_SHARDING_VERSION; import static ru.rt.restream.reindexer.binding.Consts.QUERY_RESULT_SHARD_ID; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2; import static ru.rt.restream.reindexer.binding.Consts.RANK_FORMAT_SINGLE_FLOAT; import static ru.rt.restream.reindexer.binding.Consts.RESULTS_FORMAT_MASK; import static ru.rt.restream.reindexer.binding.Consts.RESULTS_JSON; @@ -59,8 +60,27 @@ public class QueryResultReader { * @return the {@link QueryResult} to use */ public QueryResult read(byte[] rawQueryResult) { + return read(rawQueryResult, Consts.QUERY_FORMAT_V1); + } + + /** + * Reads a {@link QueryResult} from the raw byte array. + * + * @param rawQueryResult the raw byte array + * @param queryFormatVersion query serialization format version + * @return the {@link QueryResult} to use + */ + public QueryResult read(byte[] rawQueryResult, int queryFormatVersion) { ByteBuffer buffer = new ByteBuffer(rawQueryResult).rewind(); + if (queryFormatVersion == QUERY_FORMAT_V2) { + long format = buffer.getVarUInt(); + if (format != QUERY_FORMAT_V2) { + String errorMessage = String.format("QueryResults format version='%d' is not supported", format); + throw new RuntimeException(errorMessage); + } + } QueryResult queryResult = getQueryResultWithFlags(buffer.getVarUInt()); + queryResult.setQueryFormatVersion(queryFormatVersion); queryResult.setTotalCount(buffer.getVarUInt()); queryResult.setQCount(buffer.getVarUInt()); queryResult.setCount(buffer.getVarUInt()); diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java index 21655239..5cbaa66e 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java @@ -35,6 +35,7 @@ import java.time.Duration; import java.util.concurrent.atomic.AtomicLong; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2; import static ru.rt.restream.reindexer.binding.Consts.REINDEXER_VERSION; /** @@ -207,6 +208,11 @@ private void checkResponse(ReindexerResponse response) { } } + @Override + public int queryFormatVersion() { + return QUERY_FORMAT_V2; + } + @Override public void close() { adapter.destroy(rx); diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java index 49bc1965..6bc7cc24 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java @@ -22,6 +22,8 @@ import ru.rt.restream.reindexer.binding.RequestContext; import ru.rt.restream.reindexer.util.NativeUtils; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2; + /** * A request context which is holds a {@link QueryResult}, * the {@link #fetchResults(int, int)} method is NOOP since Builtin does not support it. @@ -52,7 +54,7 @@ public BuiltinRequestContext(ReindexerResponse response) { } } QueryResultReader reader = new QueryResultReader(); - queryResult = reader.read(rawQueryResult); + queryResult = reader.read(rawQueryResult, QUERY_FORMAT_V2); queryResult.setResultsPtr(resultsPtr); } diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java index bc1f8741..f83cdde1 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java @@ -191,6 +191,11 @@ public String getMeta(String namespace, String key) { return builtin.getMeta(namespace, key); } + @Override + public int queryFormatVersion() { + return builtin.queryFormatVersion(); + } + @Override public void close() { ReindexerResponse response = adapter.stopServer(svc); diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java index ed8b7fa5..1968d95d 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java @@ -50,6 +50,15 @@ public interface Connection extends AutoCloseable { */ boolean hasError(); + /** + * Returns negotiated query serialization format version. + * + * @return query serialization format version + */ + default int queryFormatVersion() { + return ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V1; + } + /** * Closes the connection. */ diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java index a394f0ff..984dc842 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java @@ -145,6 +145,10 @@ public Connection getConnection() { return connection; } + public int queryFormatVersion() { + return getConnection().queryFormatVersion(); + } + private DataSource getDataSource(int connectionPoolSize) { Instant connectionDeadline = Instant.now().plus(timeout); for (; ; ) { diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java index e369b696..531858bf 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java @@ -142,7 +142,7 @@ public RequestContext select(String query, boolean asJson, int fetchCount, long[ Connection connection = pool.getConnection(); ReindexerResponse rpcResponse = ConnectionUtils.rpcCall(connection, SELECT_SQL, query, flags, fetchCount > 0 ? fetchCount : Integer.MAX_VALUE, ptVersions); - return new CprotoRequestContext(rpcResponse, connection, asJson); + return new CprotoRequestContext(rpcResponse, connection, asJson, connection.queryFormatVersion()); } /** @@ -156,7 +156,7 @@ public RequestContext selectQuery(byte[] queryData, int fetchCount, long[] ptVer Connection connection = pool.getConnection(); ReindexerResponse rpcResponse = ConnectionUtils.rpcCall(connection, SELECT, queryData, flags, fetchCount > 0 ? fetchCount : Integer.MAX_VALUE, ptVersions); - return new CprotoRequestContext(rpcResponse, connection, asJson); + return new CprotoRequestContext(rpcResponse, connection, asJson, connection.queryFormatVersion()); } @Override @@ -191,6 +191,11 @@ public String getMeta(String namespace, String key) { return new String((byte[]) response.getArguments()[0], StandardCharsets.UTF_8); } + @Override + public int queryFormatVersion() { + return pool.queryFormatVersion(); + } + /** * Closes the connection pool. */ diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java index dd6196e7..4a4df39f 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java @@ -46,6 +46,8 @@ public class CprotoRequestContext implements RequestContext { private final boolean asJson; + private final int queryFormatVersion; + private QueryResult queryResult; private int requestId = -1; @@ -57,10 +59,12 @@ public class CprotoRequestContext implements RequestContext { * @param connection the connection in which the request was made * @param asJson 'true' if response should be serialized in JSON format, defaults to CJSON */ - public CprotoRequestContext(ReindexerResponse rpcResponse, Connection connection, boolean asJson) { - this.queryResult = getQueryResult(rpcResponse); + public CprotoRequestContext(ReindexerResponse rpcResponse, Connection connection, boolean asJson, + int queryFormatVersion) { this.connection = connection; this.asJson = asJson; + this.queryFormatVersion = queryFormatVersion; + this.queryResult = getQueryResult(rpcResponse); } @Override @@ -101,7 +105,7 @@ private QueryResult getQueryResult(ReindexerResponse rpcResponse) { if (responseArguments.length > 1) { requestId = (int) responseArguments[1]; } - return reader.read(rawQueryResult); + return reader.read(rawQueryResult, queryFormatVersion); } } diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoTransactionContext.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoTransactionContext.java index f38a06ae..faa4fb4b 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoTransactionContext.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoTransactionContext.java @@ -103,7 +103,7 @@ public RequestContext selectQuery(byte[] queryData, int fetchCount, long[] ptVer : Consts.RESULTS_C_JSON | Consts.RESULTS_WITH_PAYLOAD_TYPES; ReindexerResponse rpcResponse = ConnectionUtils.rpcCall(connection, SELECT, queryData, flags, fetchCount > 0 ? fetchCount : Integer.MAX_VALUE, ptVersions); - return new CprotoRequestContext(rpcResponse, connection, asJson); + return new CprotoRequestContext(rpcResponse, connection, asJson, connection.queryFormatVersion()); } @Override diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java index e6ad5bdd..43dfe8d1 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java @@ -52,8 +52,11 @@ import static ru.rt.restream.reindexer.binding.Consts.APP_PROPERTY_NAME; import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_COMPLEX_RANK; import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_NAMESPACE_INCARNATIONS; +import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_QUERY_FORMAT_V2; import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS; import static ru.rt.restream.reindexer.binding.Consts.DEF_APP_NAME; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V1; +import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2; import static ru.rt.restream.reindexer.binding.Consts.REINDEXER_VERSION; /** @@ -72,7 +75,7 @@ public class PhysicalConnection implements Connection { static final long CPROTO_MAGIC = 0xEEDD1132L; - static final int CPROTO_VERSION = 0x104; + static final int CPROTO_VERSION = 0x105; static final int CPROTO_HDR_LEN = 16; @@ -108,6 +111,8 @@ public class PhysicalConnection implements Connection { private final ScheduledFuture writeTaskFuture; + private volatile int queryFormatVersion = QUERY_FORMAT_V1; + public PhysicalConnection(String host, int port, String user, String password, String database, SSLSocketFactory sslSocketFactory, Duration requestTimeout, ScheduledExecutorService scheduler) { @@ -131,7 +136,7 @@ public PhysicalConnection(String host, int port, String user, String password, S } readTaskFuture = scheduler.scheduleWithFixedDelay(new ReadTask(), 0, 100, TimeUnit.MICROSECONDS); writeTaskFuture = scheduler.scheduleWithFixedDelay(new WriteTask(), 0, 100, TimeUnit.MICROSECONDS); - ConnectionUtils.rpcCallNoResults(this, Binding.LOGIN, user, password, database, + ReindexerResponse response = ConnectionUtils.rpcCall(this, Binding.LOGIN, user, password, database, false, // create DB if missing false, // checkClusterID -1, // expectedClusterID @@ -139,13 +144,25 @@ public PhysicalConnection(String host, int port, String user, String password, S getAppName(), BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS | BINDING_CAPABILITY_COMPLEX_RANK - | BINDING_CAPABILITY_NAMESPACE_INCARNATIONS); + | BINDING_CAPABILITY_NAMESPACE_INCARNATIONS + | BINDING_CAPABILITY_QUERY_FORMAT_V2); + updateQueryFormatVersion(response); } catch (Exception e) { onError(e); throw new NetworkException(e); } } + private void updateQueryFormatVersion(ReindexerResponse response) { + Object[] arguments = response.getArguments(); + if (arguments.length > 2 && arguments[2] instanceof Long) { + long capabilities = (Long) arguments[2]; + queryFormatVersion = (capabilities & BINDING_CAPABILITY_QUERY_FORMAT_V2) != 0 + ? QUERY_FORMAT_V2 + : QUERY_FORMAT_V1; + } + } + private Object getAppName() { return System.getProperty(APP_PROPERTY_NAME, DEF_APP_NAME); } @@ -335,6 +352,11 @@ public boolean hasError() { return getCurrentError() != null; } + @Override + public int queryFormatVersion() { + return queryFormatVersion; + } + private Exception getCurrentError() { lock.readLock().lock(); try {