Bump couchbase sdk to 3_4_3. (#1663)

Requires some refactoring around @Stability.Internal APIs.
Also fixed a test to get it to pass.

Closes #1661,#1662.
This commit is contained in:
Michael Reiche
2023-02-17 14:03:44 -08:00
committed by GitHub
parent ddd19744d8
commit 02afaeef8d
7 changed files with 80 additions and 51 deletions

View File

@@ -15,6 +15,8 @@
*/
package org.springframework.data.couchbase.core;
import com.couchbase.client.core.api.query.CoreQueryContext;
import com.couchbase.client.core.io.CollectionIdentifier;
import reactor.core.publisher.Flux;
import java.util.Optional;
@@ -28,7 +30,6 @@ import org.springframework.data.couchbase.core.support.PseudoArgs;
import org.springframework.data.couchbase.core.support.TemplateUtils;
import org.springframework.util.Assert;
import com.couchbase.client.core.deps.com.fasterxml.jackson.databind.node.ObjectNode;
import com.couchbase.client.java.ReactiveScope;
import com.couchbase.client.java.json.JsonObject;
import com.couchbase.client.java.query.QueryOptions;
@@ -100,11 +101,10 @@ public class ReactiveRemoveByQueryOperationSupport implements ReactiveRemoveByQu
} else {
TransactionQueryOptions opts = OptionsBuilder
.buildTransactionQueryOptions(buildQueryOptions(pArgs.getOptions()));
ObjectNode convertedOptions = com.couchbase.client.java.transactions.internal.OptionsUtil
.createTransactionOptions(pArgs.getScope() == null ? null : rs, statement, opts);
CoreQueryContext queryContext = CollectionIdentifier.DEFAULT_SCOPE.equals(rs.name()) ? null : CoreQueryContext.of(rs.bucketName(), rs.name());
return transactionContext.get().getCore()
.queryBlocking(statement, template.getBucketName(), pArgs.getScope(), convertedOptions, false)
.flatMapIterable(result -> result.rows).map(row -> {
.queryBlocking(statement, queryContext, opts.builder().build(), false)
.flatMapIterable(result -> result.collectRows()).map(row -> {
JsonObject json = JsonObject.fromJson(row.data());
return new RemoveResult(json.getString(TemplateUtils.SELECT_ID), json.getLong(TemplateUtils.SELECT_CAS),
Optional.empty());

View File

@@ -41,9 +41,10 @@ public class N1QLQuery extends Query {
return options;
}
// for logging only
public JsonObject n1ql() {
JsonObject query = JsonObject.create().put("statement", expression.toString());
options.build().injectParams(query);
query.put("options", OptionsBuilder.getQueryOpts(options.build()));
return query;
}

View File

@@ -15,6 +15,7 @@
*/
package org.springframework.data.couchbase.core.query;
import static com.couchbase.client.core.util.Validators.notNull;
import static org.springframework.data.couchbase.core.query.Meta.MetaKey.RETRY_STRATEGY;
import static org.springframework.data.couchbase.core.query.Meta.MetaKey.SCAN_CONSISTENCY;
import static org.springframework.data.couchbase.core.query.Meta.MetaKey.TIMEOUT;
@@ -24,9 +25,13 @@ import java.lang.reflect.AnnotatedElement;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.time.Duration;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import com.couchbase.client.core.api.query.CoreQueryScanConsistency;
import com.couchbase.client.core.classic.query.ClassicCoreQueryOps;
import com.couchbase.client.core.error.InvalidArgumentException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.core.annotation.AnnotatedElementUtils;
@@ -67,21 +72,21 @@ public class OptionsBuilder {
static QueryOptions buildQueryOptions(Query query, QueryOptions options, QueryScanConsistency scanConsistency) {
options = options != null ? options : QueryOptions.queryOptions();
if (query.getParameters() != null) {
if (query.getParameters() instanceof JsonArray) {
if (query.getParameters() instanceof JsonArray && !((JsonArray) query.getParameters()).isEmpty()) {
options.parameters((JsonArray) query.getParameters());
} else {
} else if( query.getParameters() instanceof JsonObject && !((JsonObject)query.getParameters()).isEmpty()){
options.parameters((JsonObject) query.getParameters());
}
}
Meta meta = query.getMeta() != null ? query.getMeta() : new Meta();
QueryOptions.Built optsBuilt = options.build();
JsonObject optsJson = getQueryOpts(optsBuilt);
QueryScanConsistency metaQueryScanConsistency = meta.get(SCAN_CONSISTENCY) != null
? ((ScanConsistency) meta.get(SCAN_CONSISTENCY)).query()
: null;
QueryScanConsistency qsc = fromFirst(QueryScanConsistency.NOT_BOUNDED, query.getScanConsistency(),
getScanConsistency(optsJson), scanConsistency, metaQueryScanConsistency);
scanConsistency(optsBuilt), scanConsistency, metaQueryScanConsistency);
Duration timeout = fromFirst(Duration.ofSeconds(0), getTimeout(optsBuilt), meta.get(TIMEOUT));
RetryStrategy retryStrategy = fromFirst(null, getRetryStrategy(optsBuilt), meta.get(RETRY_STRATEGY));
@@ -100,6 +105,21 @@ public class OptionsBuilder {
return options;
}
private static QueryScanConsistency scanConsistency(QueryOptions.Built optsBuilt){
CoreQueryScanConsistency scanConsistency = optsBuilt.scanConsistency();
if (scanConsistency == null){
return null;
}
switch (scanConsistency) {
case NOT_BOUNDED:
return QueryScanConsistency.NOT_BOUNDED;
case REQUEST_PLUS:
return QueryScanConsistency.REQUEST_PLUS;
default:
throw new InvalidArgumentException("Unknown scan consistency type " + scanConsistency, null, null);
}
}
public static TransactionQueryOptions buildTransactionQueryOptions(QueryOptions options) {
QueryOptions.Built built = options.build();
TransactionQueryOptions txOptions = TransactionQueryOptions.queryOptions();
@@ -110,8 +130,21 @@ public class OptionsBuilder {
throw new IllegalArgumentException("QueryOptions.flexIndex is not supported in a transaction");
}
Object value = optsJson.get("args");
if(value instanceof JsonObject){
txOptions.parameters((JsonObject)value);
}else if(value instanceof JsonArray) {
txOptions.parameters((JsonArray) value);
} else if(value != null) {
throw InvalidArgumentException.fromMessage(
"non-null args property was neither JsonObject(namedParameters) nor JsonArray(positionalParameters) "
+ value);
}
for (Map.Entry<String, Object> entry : optsJson.toMap().entrySet()) {
txOptions.raw(entry.getKey(), entry.getValue());
if(!entry.getKey().equals("args")) {
txOptions.raw(entry.getKey(), entry.getValue());
}
}
if (LOG.isDebugEnabled()) {
@@ -370,10 +403,8 @@ public class OptionsBuilder {
return s.toString();
}
private static JsonObject getQueryOpts(QueryOptions.Built optsBuilt) {
JsonObject jo = JsonObject.create();
optsBuilt.injectParams(jo);
return jo;
public static JsonObject getQueryOpts(QueryOptions.Built optsBuilt) {
return JsonObject.fromJson(ClassicCoreQueryOps.convertOptions(optsBuilt).toString().getBytes());
}
/**
@@ -396,18 +427,6 @@ public class OptionsBuilder {
return chosen;
}
private static QueryScanConsistency getScanConsistency(JsonObject opts) {
String str = opts.getString("scan_consistency");
if ("at_plus".equals(str)) {
return null;
}
return str == null ? null : QueryScanConsistency.valueOf(str.toUpperCase());
}
private static JsonObject getScanVectors(JsonObject opts) {
return opts.getObject("scan_vectors");
}
private static Duration getTimeout(QueryOptions.Built optsBuilt) {
Optional<Duration> timeout = optsBuilt.timeout();
return timeout.isPresent() ? timeout.get() : null;