diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseMetadata.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseMetadata.java new file mode 100644 index 000000000000..7ff193e03e03 --- /dev/null +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseMetadata.java @@ -0,0 +1,76 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.spanner; + +import com.google.spanner.v1.TransactionOptions.IsolationLevel; +import com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode; +import java.util.Objects; + +/** + * Internal container for dynamic database-level defaults queried from + * INFORMATION_SCHEMA.DATABASE_OPTIONS. Holds the database dialect, default transaction isolation + * level, and default read lock mode. + */ +final class DatabaseMetadata { + private final Dialect dialect; + private final IsolationLevel isolationLevel; + private final ReadLockMode readLockMode; + + DatabaseMetadata(Dialect dialect, IsolationLevel isolationLevel, ReadLockMode readLockMode) { + this.dialect = Objects.requireNonNull(dialect); + this.isolationLevel = Objects.requireNonNull(isolationLevel); + this.readLockMode = Objects.requireNonNull(readLockMode); + } + + Dialect getDialect() { + return dialect; + } + + IsolationLevel getIsolationLevel() { + return isolationLevel; + } + + ReadLockMode getReadLockMode() { + return readLockMode; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + DatabaseMetadata that = (DatabaseMetadata) o; + return dialect == that.dialect + && isolationLevel == that.isolationLevel + && readLockMode == that.readLockMode; + } + + @Override + public int hashCode() { + return Objects.hash(dialect, isolationLevel, readLockMode); + } + + @Override + public String toString() { + return String.format( + "DatabaseMetadata{dialect=%s, isolation=%s, lockMode=%s}", + dialect, isolationLevel, readLockMode); + } +} diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java index ece6d862f87b..94cd0dad212c 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java @@ -31,6 +31,8 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; import com.google.spanner.v1.BatchWriteResponse; +import com.google.spanner.v1.TransactionOptions.IsolationLevel; +import com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode; import java.time.Clock; import java.time.Duration; import java.time.Instant; @@ -63,14 +65,30 @@ final class MultiplexedSessionDatabaseClient extends AbstractMultiplexedSessionD */ private static final int MAX_INITIAL_CREATE_SESSION_ATTEMPTS = 10; + /** + * Statement used to query database-level default options from + * INFORMATION_SCHEMA.DATABASE_OPTIONS. This retrieves 'default_transaction_isolation', + * 'default_read_lock_mode', and 'database_dialect' so that the client can correctly configure + * transaction routing (e.g. Leader-Aware Routing) and isolation level behaviors without relying + * solely on client-side hardcoded defaults. + */ @VisibleForTesting - static final Statement DETERMINE_DIALECT_STATEMENT = + static final Statement DETERMINE_METADATA_STATEMENT = Statement.newBuilder( - "select option_value " - + "from information_schema.database_options " - + "where option_name='database_dialect'") + "SELECT OPTION_NAME, OPTION_VALUE " + + "FROM INFORMATION_SCHEMA.DATABASE_OPTIONS " + + "WHERE OPTION_NAME IN ('default_transaction_isolation', " + + "'default_read_lock_mode', 'database_dialect')") .build(); + static final String OPTION_DATABASE_DIALECT = "database_dialect"; + static final String OPTION_DEFAULT_TRANSACTION_ISOLATION = "default_transaction_isolation"; + static final String OPTION_DEFAULT_READ_LOCK_MODE = "default_read_lock_mode"; + + static final String ISOLATION_LEVEL_REPEATABLE_READ = "repeatable read"; + static final String READ_LOCK_MODE_OPTIMISTIC = "optimistic"; + static final String READ_LOCK_MODE_PESSIMISTIC = "pessimistic"; + /** * Represents a single transaction on a multiplexed session. This can be both a single-use or * multi-use transaction, and both read/write or read-only transaction. This can be compared to a @@ -274,13 +292,8 @@ public void onSessionReady(SessionImpl session) { // only start the maintainer if we actually managed to create a session in the first // place. maintainer.start(); - if (sessionClient - .getSpanner() - .getOptions() - .getSessionPoolOptions() - .isAutoDetectDialect()) { - MAINTAINER_SERVICE.submit(() -> getDialect()); - } + MAINTAINER_SERVICE.submit( + () -> session.getSessionReference().setDatabaseMetadata(getDatabaseMetadata())); } @Override @@ -371,6 +384,12 @@ AtomicLong getNumSessionsReleased() { return this.numSessionsReleased; } + @VisibleForTesting + void resetAcquiredAndReleasedCounts() { + this.numSessionsAcquired.set(0L); + this.numSessionsReleased.set(0L); + } + void close() { boolean releaseChannelUsage = false; synchronized (this) { @@ -473,32 +492,71 @@ private int getSingleUseChannelHint() { } } - private final AbstractLazyInitializer dialectSupplier = - new AbstractLazyInitializer() { + private static IsolationLevel parseIsolationLevel(String value) { + return ISOLATION_LEVEL_REPEATABLE_READ.equalsIgnoreCase(value) + ? IsolationLevel.REPEATABLE_READ + : IsolationLevel.SERIALIZABLE; + } + + private static ReadLockMode parseReadLockMode(String value) { + if (READ_LOCK_MODE_OPTIMISTIC.equalsIgnoreCase(value)) { + return ReadLockMode.OPTIMISTIC; + } else if (READ_LOCK_MODE_PESSIMISTIC.equalsIgnoreCase(value)) { + return ReadLockMode.PESSIMISTIC; + } + return ReadLockMode.READ_LOCK_MODE_UNSPECIFIED; + } + + /** + * Lazily initializes and caches {@link DatabaseMetadata} (dialect, default isolation level, and + * read lock mode). Introspects the database options once and attaches the resolved metadata to + * the current multiplexed {@link SessionReference} so subsequent transactions can resolve their + * effective modes. + */ + private final AbstractLazyInitializer metadataSupplier = + new AbstractLazyInitializer() { @Override - protected Dialect initialize() { - try (ResultSet dialectResultSet = singleUse().executeQuery(DETERMINE_DIALECT_STATEMENT)) { - if (dialectResultSet.next()) { - return Dialect.fromName(dialectResultSet.getString(0)); + protected DatabaseMetadata initialize() { + Dialect dialect = Dialect.GOOGLE_STANDARD_SQL; + IsolationLevel isolationLevel = IsolationLevel.SERIALIZABLE; + ReadLockMode readLockMode = ReadLockMode.READ_LOCK_MODE_UNSPECIFIED; + + numSessionsAcquired.decrementAndGet(); + try (ResultSet resultSet = singleUse().executeQuery(DETERMINE_METADATA_STATEMENT)) { + while (resultSet.next()) { + String name = resultSet.getString(0); + String value = resultSet.getString(1); + if (OPTION_DATABASE_DIALECT.equalsIgnoreCase(name)) { + dialect = Dialect.fromName(value); + } else if (OPTION_DEFAULT_TRANSACTION_ISOLATION.equalsIgnoreCase(name)) { + isolationLevel = parseIsolationLevel(value); + } else if (OPTION_DEFAULT_READ_LOCK_MODE.equalsIgnoreCase(name)) { + readLockMode = parseReadLockMode(value); + } } + } finally { + numSessionsReleased.decrementAndGet(); } - // This should not really happen, but it is the safest fallback value. - return Dialect.GOOGLE_STANDARD_SQL; + return new DatabaseMetadata(dialect, isolationLevel, readLockMode); } }; - @Override - public Dialect getDialect() { + DatabaseMetadata getDatabaseMetadata() { try { - return dialectSupplier.get(); + return metadataSupplier.get(); } catch (Exception exception) { throw SpannerExceptionFactory.asSpannerException(exception); } } + @Override + public Dialect getDialect() { + return getDatabaseMetadata().getDialect(); + } + Future getDialectAsync() { try { - return MAINTAINER_SERVICE.submit(dialectSupplier::get); + return MAINTAINER_SERVICE.submit(() -> getDialect()); } catch (Exception exception) { throw SpannerExceptionFactory.asSpannerException(exception); } @@ -659,8 +717,10 @@ void maintain() { new SessionConsumer() { @Override public void onSessionReady(SessionImpl session) { - multiplexedSessionReference.set( - ApiFutures.immediateFuture(session.getSessionReference())); + SessionReference sessionRef = session.getSessionReference(); + multiplexedSessionReference.set(ApiFutures.immediateFuture(sessionRef)); + MAINTAINER_SERVICE.submit( + () -> sessionRef.setDatabaseMetadata(getDatabaseMetadata())); expirationDate.set( clock .instant() diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionReference.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionReference.java index 1fd6c303ede1..fea649562d02 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionReference.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionReference.java @@ -38,6 +38,7 @@ class SessionReference { private volatile Instant lastUseTime; @Nullable private final Instant createTime; private final boolean isMultiplexed; + private volatile DatabaseMetadata databaseMetadata; SessionReference(String name, @Nullable String databaseRole, Map options) { this.options = options; @@ -92,6 +93,15 @@ boolean getIsMultiplexed() { return isMultiplexed; } + @Nullable + DatabaseMetadata getDatabaseMetadata() { + return databaseMetadata; + } + + void setDatabaseMetadata(DatabaseMetadata databaseMetadata) { + this.databaseMetadata = databaseMetadata; + } + void markUsed(Instant instant) { lastUseTime = instant; } diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java index d6353bfcdb86..0f164b31555f 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java @@ -54,6 +54,8 @@ import com.google.spanner.v1.RollbackRequest; import com.google.spanner.v1.Transaction; import com.google.spanner.v1.TransactionOptions; +import com.google.spanner.v1.TransactionOptions.IsolationLevel; +import com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode; import com.google.spanner.v1.TransactionSelector; import java.util.ArrayList; import java.util.List; @@ -223,8 +225,19 @@ public void removeListener(Runnable listener) { private CommitResponse commitResponse; private final Clock clock; + private final boolean routeToLeader; private final Map channelHint; + private static final class TransactionMode { + private final IsolationLevel isolationLevel; + private final ReadLockMode readLockMode; + + TransactionMode(IsolationLevel isolationLevel, ReadLockMode readLockMode) { + this.isolationLevel = isolationLevel; + this.readLockMode = readLockMode; + } + } + private TransactionContextImpl(Builder builder) { super(builder); this.transactionId = builder.transactionId; @@ -238,6 +251,81 @@ private TransactionContextImpl(Builder builder) { ThreadLocalRandom.current().nextLong(Long.MAX_VALUE), session.getSpanner().getOptions().isGrpcGcpExtensionEnabled()); this.previousTransactionId = builder.previousTransactionId; + + TransactionMode effectiveMode = resolveEffectiveMode(this.options); + this.routeToLeader = !canEnableLRYW(effectiveMode); + } + + /** + * Resolves the effective isolation level and read lock mode using the following precedence: 1. + * Call-site options explicitly passed to the transaction (`callSite`). 2. Client-side static + * default transaction options configured on `SpannerOptions`. 3. Database-level defaults + * queried from `INFORMATION_SCHEMA.DATABASE_OPTIONS` (`dbDefaults`). 4. Hardcoded spanner + * defaults (`SERIALIZABLE` isolation level). + */ + private TransactionMode resolveEffectiveMode(Options callSite) { + IsolationLevel isolationLevel = callSite.isolationLevel(); + ReadLockMode readLockMode = callSite.readLockMode(); + + TransactionOptions defaultTxOptions = null; + if (session.getSpanner() != null && session.getSpanner().getOptions() != null) { + defaultTxOptions = session.getSpanner().getOptions().getDefaultTransactionOptions(); + } + if (defaultTxOptions != null) { + if (isUnspecified(isolationLevel)) { + isolationLevel = defaultTxOptions.getIsolationLevel(); + } + if (isUnspecified(readLockMode) && defaultTxOptions.hasReadWrite()) { + readLockMode = defaultTxOptions.getReadWrite().getReadLockMode(); + } + } + + DatabaseMetadata dbDefaults = null; + if (session.getSessionReference() != null) { + dbDefaults = session.getSessionReference().getDatabaseMetadata(); + } + if (dbDefaults != null) { + if (isUnspecified(isolationLevel)) { + isolationLevel = dbDefaults.getIsolationLevel(); + } + if (isUnspecified(readLockMode)) { + readLockMode = dbDefaults.getReadLockMode(); + } + } + + if (isUnspecified(isolationLevel)) { + isolationLevel = IsolationLevel.SERIALIZABLE; + } + // For REPEATABLE_READ, keep lock mode unspecified/null when not explicitly set so that + // canEnableLRYW evaluates to true for Leader-Routed Read-Your-Writes. + if (isUnspecified(readLockMode) && isolationLevel != IsolationLevel.REPEATABLE_READ) { + readLockMode = ReadLockMode.PESSIMISTIC; + } + + return new TransactionMode(isolationLevel, readLockMode); + } + + private static boolean isUnspecified(IsolationLevel level) { + return level == null + || level == IsolationLevel.ISOLATION_LEVEL_UNSPECIFIED + || level == IsolationLevel.UNRECOGNIZED; + } + + private static boolean isUnspecified(ReadLockMode mode) { + return mode == null + || mode == ReadLockMode.READ_LOCK_MODE_UNSPECIFIED + || mode == ReadLockMode.UNRECOGNIZED; + } + + /** + * Determines whether Leader-Routed Read-Your-Writes (LRYW) can be enabled (`routeToLeader = + * false`). LRYW is enabled when readLockMode is OPTIMISTIC, or when isolation level is + * REPEATABLE_READ and readLockMode is unspecified. + */ + private static boolean canEnableLRYW(TransactionMode mode) { + return mode.readLockMode == ReadLockMode.OPTIMISTIC + || (mode.isolationLevel == IsolationLevel.REPEATABLE_READ + && isUnspecified(mode.readLockMode)); } @Override @@ -247,7 +335,7 @@ protected boolean isReadOnly() { @Override protected boolean isRouteToLeader() { - return true; + return routeToLeader; } private void increaseAsyncOperations() { diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/ChannelUsageTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/ChannelUsageTest.java index b3d71c4925bb..8b1a550ecc1a 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/ChannelUsageTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/ChannelUsageTest.java @@ -221,6 +221,9 @@ public void testUsesAllChannels() throws Exception { try (Spanner spanner = createSpannerOptions().getService()) { DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + client.getDialect(); + allExecuteSqlChannelHints.clear(); + executeSqlChannelHints.clear(); ExecutorService executor = Executors.newFixedThreadPool(concurrentTransactions); CountDownLatch ready = new CountDownLatch(concurrentTransactions); CountDownLatch start = new CountDownLatch(1); diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MockSpannerServiceImpl.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MockSpannerServiceImpl.java index cc44ba2f3f81..11324f5c145e 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MockSpannerServiceImpl.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MockSpannerServiceImpl.java @@ -267,6 +267,7 @@ private static int parseResumeToken(ByteString resumeToken) { /** The result of a statement that is executed on a {@link MockSpannerServiceImpl}. */ public static class StatementResult { + private enum StatementResultType { RESULT_SET, UPDATE_COUNT, @@ -317,10 +318,27 @@ public static StatementResult exception(Statement statement, StatusRuntimeExcept return new StatementResult(statement, exception); } - /** Creates a result for the query that detects the dialect that is used for the database. */ + /** + * Creates a result for the query that detects the dialect that is used for the database. + * Delegates to {@link #detectMetadataResult(Dialect)} as dialect detection uses database + * metadata. + */ public static StatementResult detectDialectResult(Dialect resultDialect) { - return StatementResult.query( - MultiplexedSessionDatabaseClient.DETERMINE_DIALECT_STATEMENT, + return detectMetadataResult(resultDialect); + } + + /** Creates a result for the query that detects the database metadata. */ + public static StatementResult detectMetadataResult(Dialect resultDialect) { + return detectMetadataResult(resultDialect, "serializable", "pessimistic"); + } + + /** + * Creates a result for the query that detects the database metadata with custom isolation and + * lock mode. + */ + public static StatementResult detectMetadataResult( + Dialect resultDialect, String defaultTxnIsolation, String defaultReadLockMode) { + ResultSet.Builder builder = ResultSet.newBuilder() .setMetadata( ResultSetMetadata.newBuilder() @@ -328,19 +346,57 @@ public static StatementResult detectDialectResult(Dialect resultDialect) { StructType.newBuilder() .addFields( Field.newBuilder() - .setName("DIALECT") + .setName("option_name") + .setType(Type.newBuilder().setCode(TypeCode.STRING).build()) + .build()) + .addFields( + Field.newBuilder() + .setName("option_value") .setType(Type.newBuilder().setCode(TypeCode.STRING).build()) .build()) .build()) - .build()) - .addRows( - ListValue.newBuilder() - .addValues( - com.google.protobuf.Value.newBuilder() - .setStringValue(resultDialect.toString()) - .build()) - .build()) - .build()); + .build()); + if (resultDialect != null) { + builder.addRows( + ListValue.newBuilder() + .addValues( + com.google.protobuf.Value.newBuilder() + .setStringValue("database_dialect") + .build()) + .addValues( + com.google.protobuf.Value.newBuilder() + .setStringValue(resultDialect.toString()) + .build()) + .build()); + } + if (defaultTxnIsolation != null) { + builder.addRows( + ListValue.newBuilder() + .addValues( + com.google.protobuf.Value.newBuilder() + .setStringValue("default_transaction_isolation") + .build()) + .addValues( + com.google.protobuf.Value.newBuilder() + .setStringValue(defaultTxnIsolation) + .build()) + .build()); + } + if (defaultReadLockMode != null) { + builder.addRows( + ListValue.newBuilder() + .addValues( + com.google.protobuf.Value.newBuilder() + .setStringValue("default_read_lock_mode") + .build()) + .addValues( + com.google.protobuf.Value.newBuilder() + .setStringValue(defaultReadLockMode) + .build()) + .build()); + } + return StatementResult.query( + MultiplexedSessionDatabaseClient.DETERMINE_METADATA_STATEMENT, builder.build()); } private static class KeepLastElementDeque extends LinkedList { @@ -598,7 +654,7 @@ private static void checkStreamException( /** * Flip this switch to true if you want the {@link - * MultiplexedSessionDatabaseClient#DETERMINE_DIALECT_STATEMENT} statement to be included in the + * MultiplexedSessionDatabaseClient#DETERMINE_METADATA_STATEMENT} statement to be included in the * recorded requests on the mock server. It is ignored by default to prevent tests that do not * expect this request to suddenly start failing. */ @@ -658,7 +714,7 @@ private static void checkStreamException( private SimulatedExecutionTime streamingReadExecutionTime = NO_EXECUTION_TIME; public MockSpannerServiceImpl() { - putStatementResult(StatementResult.detectDialectResult(Dialect.GOOGLE_STANDARD_SQL)); + putStatementResult(StatementResult.detectMetadataResult(Dialect.GOOGLE_STANDARD_SQL)); } private String generateSessionName(String database) { @@ -765,7 +821,7 @@ public void setAbortProbability(double probability) { /** * Set this to true if you want the {@link - * MultiplexedSessionDatabaseClient#DETERMINE_DIALECT_STATEMENT} statement to be included in the + * MultiplexedSessionDatabaseClient#DETERMINE_METADATA_STATEMENT} statement to be included in the * recorded requests on the mock server. It is ignored by default to prevent tests that do not * expect this request to suddenly start failing. */ @@ -1060,7 +1116,14 @@ void doDeleteSession(Session session) { @Override public void executeSql(ExecuteSqlRequest request, StreamObserver responseObserver) { - maybeFreezeAndRecordRequest(request); + boolean isDetermineMetadata = + !includeDetermineDialectStatementInRequests + && request + .getSql() + .equals(MultiplexedSessionDatabaseClient.DETERMINE_METADATA_STATEMENT.getSql()); + if (!isDetermineMetadata) { + maybeFreezeAndRecordRequest(request); + } Preconditions.checkNotNull(request.getSession()); Session session = getSession(request.getSession()); if (session == null) { @@ -1069,7 +1132,10 @@ public void executeSql(ExecuteSqlRequest request, StreamObserver resp } sessionLastUsed.put(session.getName(), Instant.now()); try { - executeSqlExecutionTime.simulateExecutionTime(exceptions, stickyGlobalExceptions, freezeLock); + if (!isDetermineMetadata) { + executeSqlExecutionTime.simulateExecutionTime( + exceptions, stickyGlobalExceptions, freezeLock); + } ByteString transactionId = getTransactionId(session, request.getTransaction()); simulateAbort(session, transactionId); Statement statement = @@ -1260,10 +1326,12 @@ public void executeBatchDml( @Override public void executeStreamingSql( ExecuteSqlRequest request, StreamObserver responseObserver) { - if (includeDetermineDialectStatementInRequests - || !request - .getSql() - .equals(MultiplexedSessionDatabaseClient.DETERMINE_DIALECT_STATEMENT.getSql())) { + boolean isDetermineMetadata = + !includeDetermineDialectStatementInRequests + && request + .getSql() + .equals(MultiplexedSessionDatabaseClient.DETERMINE_METADATA_STATEMENT.getSql()); + if (!isDetermineMetadata) { maybeFreezeAndRecordRequest(request); } Preconditions.checkNotNull(request.getSession()); @@ -1298,8 +1366,10 @@ public void executeStreamingSql( break; } } - executeStreamingSqlExecutionTime.simulateExecutionTime( - exceptions, stickyGlobalExceptions, freezeLock); + if (!isDetermineMetadata) { + executeStreamingSqlExecutionTime.simulateExecutionTime( + exceptions, stickyGlobalExceptions, freezeLock); + } // Get or start transaction if (!request.getPartitionToken().isEmpty()) { List tokens = @@ -1325,7 +1395,7 @@ public void executeStreamingSql( request.getTransaction(), request.getResumeToken(), responseObserver, - getExecuteStreamingSqlExecutionTime(), + isDetermineMetadata ? NO_EXECUTION_TIME : getExecuteStreamingSqlExecutionTime(), session.getMultiplexed()); break; case UPDATE_COUNT: diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java index c4cd8ba611f0..aa211c50e5bb 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java @@ -145,12 +145,26 @@ public void testCreateSessionDeadlineExceededWithNoSessionCreateWaitTime() throw testSpanner.close(); } + private DatabaseClientImpl getClient() { + return getClient(DatabaseId.of("p", "i", "d")); + } + + private DatabaseClientImpl getClient(DatabaseId databaseId) { + return getClient(spanner, databaseId); + } + + private DatabaseClientImpl getClient(Spanner spanner, DatabaseId databaseId) { + DatabaseClientImpl client = (DatabaseClientImpl) spanner.getDatabaseClient(databaseId); + client.getDialect(); + client.multiplexedSessionDatabaseClient.resetAcquiredAndReleasedCounts(); + return client; + } + @Test public void testMultiUseReadOnlyTransactionUsesSameSession() { // Execute two queries using the same transaction. Both queries should use the same // session, also when the maintainer has executed in the meantime. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); try (ReadOnlyTransaction transaction = client.readOnlyTransaction()) { try (ResultSet resultSet = transaction.executeQuery(STATEMENT)) { //noinspection StatementWithEmptyBody @@ -183,8 +197,7 @@ public void testNewTransactionUsesNewSession() { // Execute a single-use read-only transactions, then wait for the maintainer to replace the // current session, and then run another single-use read-only transaction. The two transactions // should use two different sessions. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); try (ResultSet resultSet = client.singleUse().executeQuery(STATEMENT)) { //noinspection StatementWithEmptyBody while (resultSet.next()) { @@ -215,12 +228,8 @@ public void testNewTransactionUsesNewSession() { public void testMaintainerMaintainsMultipleClients() { // Verify that the single-threaded shared executor that is used by the multiplexed client // maintains and replaces sessions from multiple clients. - DatabaseClientImpl client1 = - (DatabaseClientImpl) - spanner.getDatabaseClient(DatabaseId.of("p", "i", "d" + UUID.randomUUID())); - DatabaseClientImpl client2 = - (DatabaseClientImpl) - spanner.getDatabaseClient(DatabaseId.of("p", "i", "d" + UUID.randomUUID())); + DatabaseClientImpl client1 = getClient(DatabaseId.of("p", "i", "d" + UUID.randomUUID())); + DatabaseClientImpl client2 = getClient(DatabaseId.of("p", "i", "d" + UUID.randomUUID())); for (DatabaseClientImpl client : ImmutableList.of(client1, client2)) { try (ResultSet resultSet = client.singleUse().executeQuery(STATEMENT)) { @@ -287,8 +296,7 @@ public void testRetryWithTheSessionCreationWaitTime() { .build() .getService(); - DatabaseClientImpl client = - (DatabaseClientImpl) testSpanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(testSpanner, DatabaseId.of("p", "i", "d")); try (ResultSet resultSet = client.singleUse().executeQuery(STATEMENT)) { //noinspection StatementWithEmptyBody @@ -492,8 +500,7 @@ public void testUnimplementedErrorOnCreationIsPropagated() { @Test public void testWriteAtLeastOnceAborted() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -515,8 +522,7 @@ public void testWriteAtLeastOnceAborted() { @Test public void testWriteAtLeastOnce() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); Timestamp timestamp = MockSpannerTestActions.writeAtLeastOnceInsertMutation(client); assertNotNull(timestamp); @@ -537,8 +543,7 @@ public void testWriteAtLeastOnce() { @Test public void testWriteAtLeastOnceWithCommitStats() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); CommitResponse response = client.writeAtLeastOnceWithOptions( Collections.singletonList( @@ -565,8 +570,7 @@ public void testWriteAtLeastOnceWithCommitStats() { @Test public void testWriteAtLeastOnceWithOptions() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); MockSpannerTestActions.writeAtLeastOnceWithOptionsInsertMutation( client, Options.priority(RpcPriority.LOW)); @@ -587,8 +591,7 @@ public void testWriteAtLeastOnceWithOptions() { @Test public void testWriteAtLeastOnceWithTagOptions() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); MockSpannerTestActions.writeAtLeastOnceWithOptionsInsertMutation( client, Options.tag("app=spanner,env=test")); @@ -610,8 +613,7 @@ public void testWriteAtLeastOnceWithTagOptions() { @Test public void testWriteAtLeastOnceWithExcludeTxnFromChangeStreams() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); MockSpannerTestActions.writeAtLeastOnceWithOptionsInsertMutation( client, Options.excludeTxnFromChangeStreams()); @@ -634,8 +636,7 @@ public void testReadWriteTransactionUsingTransactionRunner() { // session. // During a retry (due to an ABORTED error), the transaction should use the same multiplexed // session as before, assuming the maintainer hasn't run in the meantime. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -676,8 +677,7 @@ public void testReadWriteTransactionUsingTransactionManager() { // session. // During a retry (due to an ABORTED error), the transaction should use the same multiplexed // session as before, assuming the maintainer hasn't run in the meantime. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -720,8 +720,7 @@ public void testReadWriteTransactionUsingTransactionManager() { @Test public void testMutationUsingWrite() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -756,8 +755,7 @@ public void testMutationUsingWrite() { @Test public void testMutationUsingWriteWithOptions() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); CommitResponse response = client.writeWithOptions( Collections.singletonList( @@ -786,8 +784,7 @@ public void testReadWriteTransactionUsingAsyncTransactionManager() throws Except // session as before, assuming the maintainer hasn't run in the meantime. final AtomicInteger attempt = new AtomicInteger(); CountDownLatch abortedLatch = new CountDownLatch(1); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); try (AsyncTransactionManager manager = client.transactionManagerAsync()) { TransactionContextFuture transactionContextFuture = manager.beginAsync(); while (true) { @@ -840,8 +837,7 @@ public void testReadWriteTransactionUsingAsyncRunner() throws Exception { // During a retry (due to an ABORTED error), the transaction should use the same multiplexed // session as before, assuming the maintainer hasn't run in the meantime. final AtomicInteger attempt = new AtomicInteger(); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); AsyncRunner runner = client.runAsync(); ApiFuture updateCount = runner.runAsync( @@ -905,8 +901,7 @@ public void testAsyncRunnerIsNonBlockingWithMultiplexedSession() throws Exceptio @Test public void testAbortedReadWriteTxnUsesPreviousTxnIdOnRetryWithInlineBegin() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -991,8 +986,7 @@ public void testAbortedReadWriteTxnUsesPreviousTxnIdOnRetryWithInlineBegin() { @Test public void testAbortedReadWriteTxnUsesPreviousTxnIdOnRetryWithExplicitBegin() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -1080,8 +1074,7 @@ public void testAbortedReadWriteTxnUsesPreviousTxnIdOnRetryWithExplicitBegin() { public void testPrecommitTokenForResultSet() { // This test verifies that the precommit token received from the ResultSet is properly tracked // and set in the CommitRequest. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); Long count = client @@ -1116,8 +1109,7 @@ public void testPrecommitTokenForResultSet() { public void testPrecommitTokenForExecuteBatchDmlResponse() { // This test verifies that the precommit token received from the ExecuteBatchDmlResponse is // properly tracked and set in the CommitRequest. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); long[] count = client @@ -1152,8 +1144,7 @@ public void testPrecommitTokenForExecuteBatchDmlResponse() { public void testPrecommitTokenForPartialResultSet() { // This test verifies that the precommit token received from the PartialResultSet is properly // tracked and set in the CommitRequest. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); client .readWriteTransaction() @@ -1188,8 +1179,7 @@ public void testPrecommitTokenForPartialResultSet() { public void testTxnTracksPrecommitTokenWithLatestSeqNo() { // This test ensures that the read-write transaction tracks the precommit token with the // highest sequence number and sets it in the CommitRequest. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); client .readWriteTransaction() @@ -1240,8 +1230,7 @@ public void testPrecommitTokenForTransactionResponse() { // and applied in the CommitRequest. The Transaction response includes a precommit token // only when the read-write transaction consists solely of mutations. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); client .readWriteTransaction() @@ -1275,8 +1264,7 @@ public void testPrecommitTokenForTransactionResponse() { public void testMutationOnlyCaseAborted() { // This test verifies that in the case of mutations-only, when a transaction is retried after an // ABORT, the mutation key is correctly set in the BeginTransaction request. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. @@ -1319,8 +1307,7 @@ public void testMutationOnlyCaseAborted() { @Test public void testMutationOnlyUsingTransactionManager() { // Test verifies mutation-only case within a R/W transaction via TransactionManager. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); try (TransactionManager manager = client.transactionManager()) { TransactionContext transaction = manager.begin(); @@ -1360,8 +1347,7 @@ public void testMutationOnlyUsingTransactionManager() { @Test public void testMutationOnlyUsingAsyncRunner() { // Test verifies mutation-only case within a R/W transaction via AsyncRunner. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); MockSpannerTestActions.asyncRunnerCommit(client, MoreExecutors.directExecutor()); // Verify that the mutation key is set in BeginTransactionRequest List beginTransactions = @@ -1384,8 +1370,7 @@ public void testMutationOnlyUsingAsyncRunner() { @Test public void testMutationOnlyUsingAsyncTransactionManager() { // Test verifies mutation-only case within a R/W transaction via AsyncTransactionManager. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); MockSpannerTestActions.transactionManagerAsyncCommit(client, MoreExecutors.directExecutor()); // Verify that the mutation key is set in BeginTransactionRequest @@ -1459,8 +1444,7 @@ public void testMutationOnlyCaseAbortedDuringBeginTransaction() { SimulatedExecutionTime.ofException( mockSpanner.createAbortedException(ByteString.copyFromUtf8("test")))); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); client .readWriteTransaction() @@ -1498,8 +1482,7 @@ public void testMutationOnlyUsingTransactionManagerAbortedDuringBeginTransaction SimulatedExecutionTime.ofException( mockSpanner.createAbortedException(ByteString.copyFromUtf8("test")))); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); try (TransactionManager manager = client.transactionManager()) { TransactionContext transaction = manager.begin(); @@ -1544,8 +1527,7 @@ public void testMutationOnlyUsingAsyncRunnerAbortedDuringBeginTransaction() { SimulatedExecutionTime.ofException( mockSpanner.createAbortedException(ByteString.copyFromUtf8("test")))); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); AsyncRunner runner = client.runAsync(); get( @@ -1584,8 +1566,7 @@ public void testMutationOnlyUsingTransactionManagerAsyncAbortedDuringBeginTransa SimulatedExecutionTime.ofException( mockSpanner.createAbortedException(ByteString.copyFromUtf8("test")))); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); try (AsyncTransactionManager manager = client.transactionManagerAsync()) { TransactionContextFuture transaction = manager.beginAsync(); @@ -1634,8 +1615,7 @@ public void testOtherUnimplementedError_ReadWriteTransactionStillUsesMultiplexed .withDescription("Multiplexed sessions are not supported.") .asRuntimeException())); - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Try to execute a query using single use transaction. try (ResultSet resultSet = client.singleUse().executeQuery(STATEMENT)) { @@ -1678,8 +1658,7 @@ public void testOtherUnimplementedError_ReadWriteTransactionStillUsesMultiplexed public void testReadWriteTransactionWithCommitRetryProtocolExtensionSet() { // This test simulates the commit retry protocol extension which occurs when a read-write // transaction contains read/query + mutation operations. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); client .readWriteTransaction() @@ -1736,8 +1715,7 @@ public void testReadWriteTransactionWithCommitRetryProtocolExtensionSet() { @Test public void testBatchWriteAtLeastOnce() { - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); Iterable MUTATION_GROUPS = ImmutableList.of( @@ -1784,8 +1762,7 @@ public void testBatchWriteAtLeastOnce() { // 2. Passes the ABORTED exception to the begin(AbortedException) method of a new // TransactionManager, and verifies that the transaction ID from the failed transaction is sent // during the inline begin of the first request. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -1873,8 +1850,7 @@ public void testBatchWriteAtLeastOnce() { // 2. Passes the ABORTED exception to the begin(AbortedException) method of a new // TransactionManager, and verifies that the transaction ID from the failed transaction is sent // during the inline begin of the first request. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); ByteString abortedTransactionID = null; AbortedException exception = null; @@ -1971,8 +1947,7 @@ public void testBatchWriteAtLeastOnce() { // AsyncTransactionManager, and verifies that the transaction ID from the failed transaction is // sent // during the inline begin of the first request. - DatabaseClientImpl client = - (DatabaseClientImpl) spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClientImpl client = getClient(); // Force the Commit RPC to return Aborted the first time it is called. The exception is cleared // after the first call, so the retry should succeed. mockSpanner.setCommitExecutionTime( @@ -2054,4 +2029,128 @@ private void waitForSessionToBeReplaced(DatabaseClientImpl client) { Thread.yield(); } } + + @Test + public void testDatabaseMetadata_allFieldsReturned() throws Exception { + mockSpanner.putStatementResult( + StatementResult.detectMetadataResult( + Dialect.POSTGRESQL, + MultiplexedSessionDatabaseClient.ISOLATION_LEVEL_REPEATABLE_READ, + MultiplexedSessionDatabaseClient.READ_LOCK_MODE_OPTIMISTIC)); + Spanner testSpanner = + SpannerOptions.newBuilder() + .setProjectId("test-project") + .setChannelProvider(channelProvider) + .setCredentials(NoCredentials.getInstance()) + .setSessionPoolOption(SessionPoolOptions.newBuilder().setFailOnSessionLeak().build()) + .build() + .getService(); + DatabaseClientImpl client = + (DatabaseClientImpl) testSpanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + + DatabaseMetadata metadata = client.multiplexedSessionDatabaseClient.getDatabaseMetadata(); + assertEquals(Dialect.POSTGRESQL, metadata.getDialect()); + assertEquals( + com.google.spanner.v1.TransactionOptions.IsolationLevel.REPEATABLE_READ, + metadata.getIsolationLevel()); + assertEquals( + com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode.OPTIMISTIC, + metadata.getReadLockMode()); + + // Execute a query to ensure session is created and session reference has the metadata. + try (ResultSet resultSet = client.singleUse().executeQuery(STATEMENT)) { + while (resultSet.next()) {} + } + assertNotNull(client.multiplexedSessionDatabaseClient.getCurrentSessionReference()); + assertEquals( + metadata, + client.multiplexedSessionDatabaseClient.getCurrentSessionReference().getDatabaseMetadata()); + } + + @Test + public void testDatabaseMetadata_missingFieldsFallbackToDefaults() throws Exception { + // Only return dialect, omit transaction isolation and read lock mode. + mockSpanner.putStatementResult( + StatementResult.detectMetadataResult(Dialect.GOOGLE_STANDARD_SQL, null, null)); + Spanner testSpanner1 = + SpannerOptions.newBuilder() + .setProjectId("test-project") + .setChannelProvider(channelProvider) + .setCredentials(NoCredentials.getInstance()) + .setSessionPoolOption(SessionPoolOptions.newBuilder().setFailOnSessionLeak().build()) + .build() + .getService(); + DatabaseClientImpl client1 = + (DatabaseClientImpl) testSpanner1.getDatabaseClient(DatabaseId.of("p", "i", "d1")); + + DatabaseMetadata metadata1 = client1.multiplexedSessionDatabaseClient.getDatabaseMetadata(); + assertEquals(Dialect.GOOGLE_STANDARD_SQL, metadata1.getDialect()); + assertEquals( + com.google.spanner.v1.TransactionOptions.IsolationLevel.SERIALIZABLE, + metadata1.getIsolationLevel()); + assertEquals( + com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode.READ_LOCK_MODE_UNSPECIFIED, + metadata1.getReadLockMode()); + + // Now test when all fields (including dialect) are omitted (empty result set). + mockSpanner.putStatementResult(StatementResult.detectMetadataResult(null, null, null)); + Spanner testSpanner2 = + SpannerOptions.newBuilder() + .setProjectId("test-project") + .setChannelProvider(channelProvider) + .setCredentials(NoCredentials.getInstance()) + .setSessionPoolOption(SessionPoolOptions.newBuilder().setFailOnSessionLeak().build()) + .build() + .getService(); + DatabaseClientImpl client2 = + (DatabaseClientImpl) testSpanner2.getDatabaseClient(DatabaseId.of("p", "i", "d2")); + + DatabaseMetadata metadata2 = client2.multiplexedSessionDatabaseClient.getDatabaseMetadata(); + assertEquals(Dialect.GOOGLE_STANDARD_SQL, metadata2.getDialect()); + assertEquals( + com.google.spanner.v1.TransactionOptions.IsolationLevel.SERIALIZABLE, + metadata2.getIsolationLevel()); + assertEquals( + com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode.READ_LOCK_MODE_UNSPECIFIED, + metadata2.getReadLockMode()); + + try (ResultSet resultSet = client2.singleUse().executeQuery(STATEMENT)) { + while (resultSet.next()) {} + } + assertNotNull(client2.multiplexedSessionDatabaseClient.getCurrentSessionReference()); + assertEquals( + metadata2, + client2 + .multiplexedSessionDatabaseClient + .getCurrentSessionReference() + .getDatabaseMetadata()); + } + + @Test + public void testDatabaseMetadata_repeatableReadDefaultLockMode() throws Exception { + mockSpanner.putStatementResult( + StatementResult.detectMetadataResult( + Dialect.GOOGLE_STANDARD_SQL, + MultiplexedSessionDatabaseClient.ISOLATION_LEVEL_REPEATABLE_READ, + null)); + Spanner testSpanner = + SpannerOptions.newBuilder() + .setProjectId("test-project") + .setChannelProvider(channelProvider) + .setCredentials(NoCredentials.getInstance()) + .setSessionPoolOption(SessionPoolOptions.newBuilder().setFailOnSessionLeak().build()) + .build() + .getService(); + DatabaseClientImpl client = + (DatabaseClientImpl) testSpanner.getDatabaseClient(DatabaseId.of("p", "i", "d-rr")); + + DatabaseMetadata metadata = client.multiplexedSessionDatabaseClient.getDatabaseMetadata(); + assertEquals(Dialect.GOOGLE_STANDARD_SQL, metadata.getDialect()); + assertEquals( + com.google.spanner.v1.TransactionOptions.IsolationLevel.REPEATABLE_READ, + metadata.getIsolationLevel()); + assertEquals( + com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode.READ_LOCK_MODE_UNSPECIFIED, + metadata.getReadLockMode()); + } } diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetryApiTracerTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetryApiTracerTest.java index 67012ed96225..f30e71dab9e7 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetryApiTracerTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetryApiTracerTest.java @@ -144,6 +144,8 @@ public void createSpannerInstance() { .build() .getService(); client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + client.getDialect(); + spanExporter.reset(); } @Test @@ -437,6 +439,8 @@ public boolean isEnableApiTracing() { .build() .getService(); DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + client.getDialect(); + spanExporter.reset(); try (ResultSet resultSet = client.singleUse().executeQuery(SELECT_RANDOM)) { assertTrue(resultSet.next()); diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetrySpanTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetrySpanTest.java index 8ff8827664da..9be4c4798671 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetrySpanTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/OpenTelemetrySpanTest.java @@ -390,6 +390,7 @@ public void transactionRunner() { "CloudSpannerOperation.CreateMultiplexedSession", "CloudSpannerOperation.ExecuteUpdate", "CloudSpannerOperation.Commit", + "CloudSpannerOperation.ExecuteStreamingQuery", "CloudSpanner.ReadWriteTransaction"); expectedReadWriteTransactionEvents = @@ -402,6 +403,7 @@ public void transactionRunner() { DatabaseClient client = getClient(); TransactionRunner runner = client.readWriteTransaction(); runner.run(transaction -> transaction.executeUpdate(UPDATE_STATEMENT)); + client.getDialect(); // Wait until the list of spans contains "CloudSpannerOperation.CreateSession", as this is // an async operation. Stopwatch stopwatch = Stopwatch.createStarted(); @@ -423,6 +425,8 @@ public void transactionRunner() { expectedCreateMultiplexedSessionsRequestEvents, expectedCreateMultiplexedSessionsRequestEventsCount); break; + case "CloudSpannerOperation.ExecuteStreamingQuery": + break; case "CloudSpannerOperation.Commit": case "CloudSpannerOperation.ExecuteUpdate": assertEquals(0, spanItem.getEvents().size()); @@ -448,6 +452,7 @@ public void transactionRunnerWithError() { ImmutableList.of( "CloudSpannerOperation.CreateMultiplexedSession", "CloudSpannerOperation.ExecuteUpdate", + "CloudSpannerOperation.ExecuteStreamingQuery", "CloudSpanner.ReadWriteTransaction"); expectedReadWriteTransactionErrorEvents = ImmutableList.of( @@ -462,6 +467,7 @@ public void transactionRunnerWithError() { SpannerException.class, () -> runner.run(transaction -> transaction.executeUpdate(INVALID_UPDATE_STATEMENT))); assertEquals(ErrorCode.INVALID_ARGUMENT, e.getErrorCode()); + client.getDialect(); List actualSpanItems = new ArrayList<>(); spanExporter @@ -483,6 +489,8 @@ public void transactionRunnerWithError() { expectedReadWriteTransactionErrorEventsCount); verifyCommonAttributes(spanItem); break; + case "CloudSpannerOperation.ExecuteStreamingQuery": + break; case "CloudSpannerOperation.ExecuteUpdate": assertEquals(0, spanItem.getEvents().size()); break; @@ -545,17 +553,17 @@ public void transactionRunnerWithFailedAndBeginTransaction() { .getFinishedSpanItems() .forEach( spanItem -> { - // Ignore multiplexed sessions, as they are not used by this test and can therefore - // best be ignored, as it is not 100% certain that it has already been created. - if (!"CloudSpannerOperation.CreateMultiplexedSession".equals(spanItem.getName())) { + // Ignore multiplexed sessions and background metadata queries, as they can best be + // ignored because it is not 100% certain when they are executed in the background. + if (!"CloudSpannerOperation.CreateMultiplexedSession".equals(spanItem.getName()) + && !"CloudSpannerOperation.ExecuteStreamingQuery".equals(spanItem.getName()) + && !"Spanner.ExecuteStreamingSql".equals(spanItem.getName())) { actualSpanItems.add(spanItem.getName()); } switch (spanItem.getName()) { case "CloudSpannerOperation.CreateMultiplexedSession": - verifyRequestEvents( - spanItem, - expectedCreateMultiplexedSessionsRequestEvents, - expectedCreateMultiplexedSessionsRequestEventsCount); + case "CloudSpannerOperation.ExecuteStreamingQuery": + case "Spanner.ExecuteStreamingSql": break; case "CloudSpannerOperation.Commit": case "CloudSpannerOperation.BeginTransaction": @@ -594,9 +602,10 @@ public void testTransactionRunnerWithRetryOnBeginTransaction() { transaction.buffer(Mutation.newInsertBuilder("foo").set("id").to(1L).build()); return null; }); + clientWithApiTracing.getDialect(); assertEquals(2, mockSpanner.countRequestsOfType(BeginTransactionRequest.class)); - int numExpectedSpans = 7; + int numExpectedSpans = 9; waitForFinishedSpans(numExpectedSpans); List finishedSpans = spanExporter.getFinishedSpanItems(); List finishedSpanNames = @@ -609,9 +618,12 @@ public void testTransactionRunnerWithRetryOnBeginTransaction() { assertTrue( actualSpanNames, finishedSpanNames.contains("CloudSpannerOperation.BeginTransaction")); assertTrue(actualSpanNames, finishedSpanNames.contains("CloudSpannerOperation.Commit")); + assertTrue( + actualSpanNames, finishedSpanNames.contains("CloudSpannerOperation.ExecuteStreamingQuery")); assertTrue(actualSpanNames, finishedSpanNames.contains("Spanner.BeginTransaction")); assertTrue(actualSpanNames, finishedSpanNames.contains("Spanner.Commit")); + assertTrue(actualSpanNames, finishedSpanNames.contains("Spanner.ExecuteStreamingSql")); SpanData beginTransactionSpan = finishedSpans.stream() @@ -638,9 +650,10 @@ public void testSingleUseRetryOnExecuteStreamingSql() { assertTrue(resultSet.next()); assertFalse(resultSet.next()); } + clientWithApiTracing.getDialect(); assertEquals(2, mockSpanner.countRequestsOfType(ExecuteSqlRequest.class)); - int numExpectedSpans = 6; + int numExpectedSpans = 8; waitForFinishedSpans(numExpectedSpans); List finishedSpans = spanExporter.getFinishedSpanItems(); List finishedSpanNames = @@ -659,8 +672,13 @@ public void testSingleUseRetryOnExecuteStreamingSql() { // means that the retry event is on this span. SpanData executeStreamingQuery = finishedSpans.stream() - .filter(span -> span.getName().equals("CloudSpannerOperation.ExecuteStreamingQuery")) - .findAny() + .filter( + span -> + span.getName().equals("CloudSpannerOperation.ExecuteStreamingQuery") + && span.getEvents().stream() + .anyMatch( + event -> event.getName().contains("Stream broken. Safe to retry"))) + .findFirst() .orElseThrow(IllegalStateException::new); assertTrue( executeStreamingQuery.toString(), @@ -681,9 +699,10 @@ public void testRetryOnExecuteSql() { clientWithApiTracing .readWriteTransaction() .run(transaction -> transaction.executeUpdate(UPDATE_STATEMENT)); + clientWithApiTracing.getDialect(); assertEquals(2, mockSpanner.countRequestsOfType(ExecuteSqlRequest.class)); - int numExpectedSpans = 7; + int numExpectedSpans = 9; waitForFinishedSpans(numExpectedSpans); List finishedSpans = spanExporter.getFinishedSpanItems(); List finishedSpanNames = @@ -694,8 +713,11 @@ public void testRetryOnExecuteSql() { assertTrue(actualSpanNames, finishedSpanNames.contains("CloudSpanner.ReadWriteTransaction")); assertTrue(actualSpanNames, finishedSpanNames.contains("CloudSpannerOperation.Commit")); + assertTrue( + actualSpanNames, finishedSpanNames.contains("CloudSpannerOperation.ExecuteStreamingQuery")); assertTrue(actualSpanNames, finishedSpanNames.contains("Spanner.ExecuteSql")); assertTrue(actualSpanNames, finishedSpanNames.contains("Spanner.Commit")); + assertTrue(actualSpanNames, finishedSpanNames.contains("Spanner.ExecuteStreamingSql")); SpanData executeSqlSpan = finishedSpans.stream() diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RequestIdMockServerTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RequestIdMockServerTest.java index eac63010915f..e97b13c915c4 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RequestIdMockServerTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RequestIdMockServerTest.java @@ -116,17 +116,39 @@ public ServerCall.Listener interceptCall( ServerCall call, Metadata headers, ServerCallHandler next) { + XGoogSpannerRequestId tmpId = null; try { String requestId = headers.get(XGoogSpannerRequestId.REQUEST_ID_HEADER_KEY); if (requestId != null) { - requestIds.add(XGoogSpannerRequestId.of(requestId)); + tmpId = XGoogSpannerRequestId.of(requestId); } else { - requestIds.add(XGoogSpannerRequestId.of(0, 0, 0, 0)); + tmpId = XGoogSpannerRequestId.of(0, 0, 0, 0); } } catch (Throwable t) { - // Ignore and continue + tmpId = null; } - return Contexts.interceptCall(Context.current(), call, headers, next); + final XGoogSpannerRequestId id = tmpId; + ServerCall.Listener listener = + Contexts.interceptCall(Context.current(), call, headers, next); + return new io.grpc.ForwardingServerCallListener + .SimpleForwardingServerCallListener< + ReqT>(listener) { + @Override + public void onMessage(ReqT message) { + boolean isDetermineMetadata = + message instanceof ExecuteSqlRequest + && ((ExecuteSqlRequest) message) + .getSql() + .equals( + MultiplexedSessionDatabaseClient + .DETERMINE_METADATA_STATEMENT + .getSql()); + if (!isDetermineMetadata && id != null) { + requestIds.add(id); + } + super.onMessage(message); + } + }; } }) .build() @@ -184,10 +206,10 @@ public static void teardown() throws InterruptedException { @Before public void prepareTest() { - // Call getClient() to make sure the multiplexed session has been created. - // Then clear all requests that were received as part of that so we don't need to include - // that in the test verifications. - getClient(); + // Call getClient().getDialect() to make sure the multiplexed session and metadata have been + // initialized. Then clear all requests that were received as part of that so we don't need to + // include that in the test verifications. + getClient().getDialect(); mockSpanner.reset(); requestIds.clear(); ((SpannerImpl) spanner).resetRequestIdCounters(); @@ -660,6 +682,7 @@ public void testOtherClientId() { try (Spanner spanner = createSpanner()) { otherClientId = ((SpannerImpl) spanner).getRequestIdClientId(); DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + client.getDialect(); try (ResultSet resultSet = client.singleUse().executeQuery(SELECT1)) { while (resultSet.next()) {} } @@ -683,7 +706,7 @@ public void testOtherClientId() { // the requests that we see. This request does not include a channel hint, hence the // zero value for the channel number in the request ID. XGoogSpannerRequestId.of(otherClientId, 0, 1, 1), - XGoogSpannerRequestId.of(otherClientId, -1, 2, 1), + XGoogSpannerRequestId.of(otherClientId, -1, 3, 1), XGoogSpannerRequestId.of(getClientId(), -1, 2, 1)), actual); } diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RetryOnDifferentGrpcChannelMockServerTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RetryOnDifferentGrpcChannelMockServerTest.java index 8dc5fdf53853..30cfef2457f8 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RetryOnDifferentGrpcChannelMockServerTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RetryOnDifferentGrpcChannelMockServerTest.java @@ -324,11 +324,12 @@ public void testSingleUseQuery_retriesOnNewChannel() { SpannerOptions.Builder builder = createSpannerOptionsBuilder(); builder.setSessionPoolOption( SessionPoolOptions.newBuilder().setUseMultiplexedSession(true).build()); - mockSpanner.setExecuteStreamingSqlExecutionTime( - SimulatedExecutionTime.ofException(Status.DEADLINE_EXCEEDED.asRuntimeException())); - try (Spanner spanner = builder.build().getService()) { DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + client.getDialect(); + ACTUAL_CHANNEL_IDS.clear(); + mockSpanner.setExecuteStreamingSqlExecutionTime( + SimulatedExecutionTime.ofException(Status.DEADLINE_EXCEEDED.asRuntimeException())); try (ResultSet resultSet = client.singleUse().executeQuery(SELECT1_STATEMENT)) { assertTrue(resultSet.next()); assertEquals(1L, resultSet.getLong(0)); diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpanTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpanTest.java index 449e78cf6125..0d77ddcc871c 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpanTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpanTest.java @@ -379,6 +379,7 @@ public void multiUse() { public void transactionRunner() { TransactionRunner runner = client.readWriteTransaction(); runner.run(transaction -> transaction.executeUpdate(UPDATE_STATEMENT)); + client.getDialect(); Map spans = failOnOverkillTraceComponent.getSpans(); assertThat(spans).containsEntry("CloudSpanner.ReadWriteTransaction", true); assertThat(spans).containsEntry("CloudSpannerOperation.Commit", true); @@ -389,6 +390,7 @@ public void transactionRunner() { "Starting Commit", "Commit Done", "Transaction Attempt Succeeded", + "Starting/Resuming stream", "Request for 1 multiplexed session returned 1 session"); verifyAnnotations( failOnOverkillTraceComponent.getAnnotations().stream() @@ -405,18 +407,20 @@ public void transactionRunnerWithError() { SpannerException.class, () -> runner.run(transaction -> transaction.executeUpdate(INVALID_UPDATE_STATEMENT))); assertEquals(ErrorCode.INVALID_ARGUMENT, e.getErrorCode()); - + client.getDialect(); Map spans = failOnOverkillTraceComponent.getSpans(); - assertEquals(spans.toString(), 3, spans.size()); + assertEquals(spans.toString(), 4, spans.size()); assertThat(spans).containsEntry("CloudSpannerOperation.CreateMultiplexedSession", true); assertThat(spans).containsEntry("CloudSpanner.ReadWriteTransaction", true); assertThat(spans).containsEntry("CloudSpannerOperation.ExecuteUpdate", true); + assertThat(spans).containsEntry("CloudSpannerOperation.ExecuteStreamingQuery", true); List expectedAnnotationsForMultiplexedSessionsRW = ImmutableList.of( "Starting Transaction Attempt", "Transaction Attempt Failed in user operation", + "Starting/Resuming stream", "Request for 1 multiplexed session returned 1 session"); verifyAnnotations( failOnOverkillTraceComponent.getAnnotations().stream() diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionChannelHintTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionChannelHintTest.java index bccc7e7da125..84186d341160 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionChannelHintTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionChannelHintTest.java @@ -257,10 +257,23 @@ private static void assertSingleRemoteClientPort(Set... remotePortSets) assertEquals(1, allRemotePorts.size()); } + private DatabaseClient getClient(Spanner spanner) { + DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + client.getDialect(); + executeSqlAffinityKeys.clear(); + beginTransactionAffinityKeys.clear(); + streamingReadAffinityKeys.clear(); + executeSqlRemotePorts.clear(); + beginTransactionRemotePorts.clear(); + streamingReadRemotePorts.clear(); + commitRemotePorts.clear(); + return client; + } + @Test public void testSingleUseReadOnlyTransaction_usesSingleChannelHint() { try (Spanner spanner = createSpannerOptions().getService()) { - DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClient client = getClient(spanner); try (ResultSet resultSet = client.singleUseReadOnlyTransaction().executeQuery(SELECT1)) { while (resultSet.next()) {} } @@ -272,7 +285,7 @@ public void testSingleUseReadOnlyTransaction_usesSingleChannelHint() { @Test public void testSingleUseReadOnlyTransaction_withTimestampBound_usesSingleChannelHint() { try (Spanner spanner = createSpannerOptions().getService()) { - DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClient client = getClient(spanner); try (ResultSet resultSet = client .singleUseReadOnlyTransaction(TimestampBound.ofExactStaleness(15L, TimeUnit.SECONDS)) @@ -287,7 +300,7 @@ public void testSingleUseReadOnlyTransaction_withTimestampBound_usesSingleChanne @Test public void testReadOnlyTransaction_usesSingleChannelHint() { try (Spanner spanner = createSpannerOptions().getService()) { - DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClient client = getClient(spanner); try (ReadOnlyTransaction transaction = client.readOnlyTransaction()) { try (ResultSet resultSet = transaction.executeQuery(SELECT1)) { while (resultSet.next()) {} @@ -307,7 +320,7 @@ public void testReadOnlyTransaction_usesSingleChannelHint() { @Test public void testReadOnlyTransaction_withTimestampBound_usesSingleChannelHint() { try (Spanner spanner = createSpannerOptions().getService()) { - DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClient client = getClient(spanner); try (ReadOnlyTransaction transaction = client.readOnlyTransaction(TimestampBound.ofExactStaleness(15L, TimeUnit.SECONDS))) { try (ResultSet resultSet = transaction.executeQuery(SELECT1)) { @@ -328,7 +341,7 @@ public void testReadOnlyTransaction_withTimestampBound_usesSingleChannelHint() { @Test public void testTransactionManager_usesSingleChannelHint() { try (Spanner spanner = createSpannerOptions().getService()) { - DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClient client = getClient(spanner); try (TransactionManager manager = client.transactionManager()) { TransactionContext transaction = manager.begin(); while (true) { @@ -359,7 +372,7 @@ public void testTransactionManager_usesSingleChannelHint() { @Test public void testTransactionRunner_usesSingleChannelHint() { try (Spanner spanner = createSpannerOptions().getService()) { - DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of("p", "i", "d")); + DatabaseClient client = getClient(spanner); TransactionRunner runner = client.readWriteTransaction(); runner.run( transaction -> { diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionContextImplTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionContextImplTest.java index 49a47364a58a..d6a818c81e38 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionContextImplTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionContextImplTest.java @@ -37,6 +37,7 @@ import com.google.spanner.v1.CommitRequest; import com.google.spanner.v1.ExecuteBatchDmlRequest; import com.google.spanner.v1.ExecuteBatchDmlResponse; +import com.google.spanner.v1.TransactionOptions; import io.opentelemetry.api.common.Attributes; import java.util.Collections; import org.junit.Before; @@ -72,6 +73,8 @@ public void setup() { when(session.getRequestIdCreator()).thenReturn(NoopRequestIdCreator.INSTANCE); SpannerImpl spanner = mock(SpannerImpl.class); SpannerOptions spannerOptions = mock(SpannerOptions.class); + when(spannerOptions.getDefaultTransactionOptions()) + .thenReturn(TransactionOptions.getDefaultInstance()); when(spanner.getOptions()).thenReturn(spannerOptions); when(session.getSpanner()).thenReturn(spanner); doNothing().when(span).setStatus(any(Throwable.class)); @@ -220,6 +223,8 @@ private void batchDml(int status) { when(session.getRequestIdCreator()).thenReturn(NoopRequestIdCreator.INSTANCE); SpannerImpl spanner = mock(SpannerImpl.class); SpannerOptions spannerOptions = mock(SpannerOptions.class); + when(spannerOptions.getDefaultTransactionOptions()) + .thenReturn(TransactionOptions.getDefaultInstance()); when(spanner.getOptions()).thenReturn(spannerOptions); when(session.getSpanner()).thenReturn(spanner); SpannerRpc rpc = mock(SpannerRpc.class); diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRoutingTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRoutingTest.java new file mode 100644 index 000000000000..b7a9357c3c4a --- /dev/null +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRoutingTest.java @@ -0,0 +1,201 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.spanner; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import com.google.cloud.spanner.TransactionRunnerImpl.TransactionContextImpl; +import com.google.cloud.spanner.spi.v1.SpannerRpc; +import com.google.protobuf.ByteString; +import com.google.spanner.v1.TransactionOptions; +import com.google.spanner.v1.TransactionOptions.IsolationLevel; +import com.google.spanner.v1.TransactionOptions.ReadWrite.ReadLockMode; +import java.util.UUID; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +@RunWith(JUnit4.class) +public class TransactionRoutingTest { + + private SpannerRpc rpc; + private SessionImpl session; + private SessionReference sessionReference; + private SpannerImpl spanner; + private SpannerOptions spannerOptions; + + @Before + public void setUp() { + rpc = mock(SpannerRpc.class); + session = mock(SessionImpl.class); + sessionReference = mock(SessionReference.class); + spanner = mock(SpannerImpl.class); + spannerOptions = mock(SpannerOptions.class); + + when(session.getSessionReference()).thenReturn(sessionReference); + when(session.getSpanner()).thenReturn(spanner); + when(spanner.getOptions()).thenReturn(spannerOptions); + + // Default: system-level unspecified (Tier 4 fallback) + setClientDefaultOptions( + IsolationLevel.ISOLATION_LEVEL_UNSPECIFIED, ReadLockMode.READ_LOCK_MODE_UNSPECIFIED); + + TraceWrapper tracer = mock(TraceWrapper.class); + ISpan span = mock(ISpan.class); + when(session.getTracer()).thenReturn(tracer); + when(tracer.getCurrentSpan()).thenReturn(span); + } + + private void setClientDefaultOptions(IsolationLevel iso, ReadLockMode lock) { + TransactionOptions defaultTxnOptions = + TransactionOptions.newBuilder() + .setIsolationLevel(iso) + .setReadWrite(TransactionOptions.ReadWrite.newBuilder().setReadLockMode(lock)) + .build(); + when(spannerOptions.getDefaultTransactionOptions()).thenReturn(defaultTxnOptions); + } + + private TransactionContextImpl createContext(Options options) { + return TransactionContextImpl.newBuilder() + .setSession(session) + .setTransactionId(ByteString.copyFromUtf8(UUID.randomUUID().toString())) + .setOptions(options) + .setRpc(rpc) + .setTracer(session.getTracer()) + .setSpan(session.getTracer().getCurrentSpan()) + .build(); + } + + @Test + public void testPrecedenceTier4SystemDefaults() { + // Tier 4: Unspecified everywhere should fall back to SERIALIZABLE / PESSIMISTIC, + // routeToLeader=true. + TransactionContextImpl context = createContext(Options.fromTransactionOptions()); + assertTrue(context.isRouteToLeader()); + } + + @Test + public void testPrecedenceTier3DatabaseDefaultsRepeatableRead() { + // Tier 3: Database Defaults sets REPEATABLE_READ + OPTIMISTIC, routeToLeader=false. + DatabaseMetadata dbMetadata = + new DatabaseMetadata( + Dialect.GOOGLE_STANDARD_SQL, IsolationLevel.REPEATABLE_READ, ReadLockMode.OPTIMISTIC); + when(sessionReference.getDatabaseMetadata()).thenReturn(dbMetadata); + + TransactionContextImpl context = createContext(Options.fromTransactionOptions()); + assertFalse(context.isRouteToLeader()); + } + + @Test + public void testPrecedenceTier3DatabaseDefaultsRepeatableReadUnspecifiedLock() { + // Tier 3: Database Defaults sets REPEATABLE_READ + unspecified lock, routeToLeader=false. + DatabaseMetadata dbMetadata = + new DatabaseMetadata( + Dialect.GOOGLE_STANDARD_SQL, + IsolationLevel.REPEATABLE_READ, + ReadLockMode.READ_LOCK_MODE_UNSPECIFIED); + when(sessionReference.getDatabaseMetadata()).thenReturn(dbMetadata); + + TransactionContextImpl context = createContext(Options.fromTransactionOptions()); + assertFalse(context.isRouteToLeader()); + } + + @Test + public void testPrecedenceTier3DatabaseDefaultsSerializable() { + // Tier 3: Database Defaults sets SERIALIZABLE + PESSIMISTIC, routeToLeader=true. + DatabaseMetadata dbMetadata = + new DatabaseMetadata( + Dialect.GOOGLE_STANDARD_SQL, IsolationLevel.SERIALIZABLE, ReadLockMode.PESSIMISTIC); + when(sessionReference.getDatabaseMetadata()).thenReturn(dbMetadata); + + TransactionContextImpl context = createContext(Options.fromTransactionOptions()); + assertTrue(context.isRouteToLeader()); + } + + @Test + public void testPrecedenceTier2ClientStaticOverridesDatabase() { + // Tier 2: Client-Side Static config is REPEATABLE_READ + OPTIMISTIC. + // Database Defaults has SERIALIZABLE + PESSIMISTIC. + // routeToLeader=false because Tier 2 overrides Tier 3. + setClientDefaultOptions(IsolationLevel.REPEATABLE_READ, ReadLockMode.OPTIMISTIC); + + DatabaseMetadata dbMetadata = + new DatabaseMetadata( + Dialect.GOOGLE_STANDARD_SQL, IsolationLevel.SERIALIZABLE, ReadLockMode.PESSIMISTIC); + when(sessionReference.getDatabaseMetadata()).thenReturn(dbMetadata); + + TransactionContextImpl context = createContext(Options.fromTransactionOptions()); + assertFalse(context.isRouteToLeader()); + } + + @Test + public void testPrecedenceTier1CallSiteOverridesAll() { + // Tier 1: Call-site options explicitly set REPEATABLE_READ + OPTIMISTIC. + // Client-Side Static has SERIALIZABLE + PESSIMISTIC. + // Database Defaults has SERIALIZABLE + PESSIMISTIC. + // routeToLeader=false because Tier 1 overrides everything. + setClientDefaultOptions(IsolationLevel.SERIALIZABLE, ReadLockMode.PESSIMISTIC); + + DatabaseMetadata dbMetadata = + new DatabaseMetadata( + Dialect.GOOGLE_STANDARD_SQL, IsolationLevel.SERIALIZABLE, ReadLockMode.PESSIMISTIC); + when(sessionReference.getDatabaseMetadata()).thenReturn(dbMetadata); + + Options callSiteOptions = + Options.fromTransactionOptions( + Options.isolationLevel(IsolationLevel.REPEATABLE_READ), + Options.readLockMode(ReadLockMode.OPTIMISTIC)); + + TransactionContextImpl context = createContext(callSiteOptions); + assertFalse(context.isRouteToLeader()); + } + + @Test + public void testPrecedenceTier1CallSiteSerializableOverridesAll() { + // Tier 1: Call-site options explicitly set SERIALIZABLE + PESSIMISTIC. + // Client-Side Static has REPEATABLE_READ + OPTIMISTIC. + // Database Defaults has REPEATABLE_READ + OPTIMISTIC. + // routeToLeader=true because Tier 1 overrides everything. + setClientDefaultOptions(IsolationLevel.REPEATABLE_READ, ReadLockMode.OPTIMISTIC); + + DatabaseMetadata dbMetadata = + new DatabaseMetadata( + Dialect.GOOGLE_STANDARD_SQL, IsolationLevel.REPEATABLE_READ, ReadLockMode.OPTIMISTIC); + when(sessionReference.getDatabaseMetadata()).thenReturn(dbMetadata); + + Options callSiteOptions = + Options.fromTransactionOptions( + Options.isolationLevel(IsolationLevel.SERIALIZABLE), + Options.readLockMode(ReadLockMode.PESSIMISTIC)); + + TransactionContextImpl context = createContext(callSiteOptions); + assertTrue(context.isRouteToLeader()); + } + + @Test + public void testRepeatableReadWithNullLockModeCanEnableLRYW() { + when(sessionReference.getDatabaseMetadata()).thenReturn(null); + Options callSiteOptions = + Options.fromTransactionOptions(Options.isolationLevel(IsolationLevel.REPEATABLE_READ)); + TransactionContextImpl context = createContext(callSiteOptions); + assertFalse(context.isRouteToLeader()); + } +} diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java index 1dd2418aa05a..b1ad6573ab3d 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java @@ -127,6 +127,8 @@ public void setUp() { when(rpc.getRequestIdCreator()).thenReturn(NoopRequestIdCreator.INSTANCE); SpannerImpl spanner = mock(SpannerImpl.class); SpannerOptions spannerOptions = mock(SpannerOptions.class); + when(spannerOptions.getDefaultTransactionOptions()) + .thenReturn(TransactionOptions.getDefaultInstance()); when(spanner.getOptions()).thenReturn(spannerOptions); when(session.getSpanner()).thenReturn(spanner); when(rpc.executeQuery(Mockito.any(ExecuteSqlRequest.class), Mockito.anyMap(), eq(true))) diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AbstractMockServerTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AbstractMockServerTest.java index 43a6b9d4feff..829adf7a697a 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AbstractMockServerTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AbstractMockServerTest.java @@ -206,7 +206,7 @@ public void getOperation( StatementResult.query(SELECT_RANDOM_STATEMENT, RANDOM_RESULT_SET)); mockSpanner.putStatementResult(StatementResult.query(SELECT1_STATEMENT, SELECT1_RESULTSET)); mockSpanner.putStatementResult( - StatementResult.detectDialectResult(Dialect.GOOGLE_STANDARD_SQL)); + MockSpannerServiceImpl.StatementResult.detectMetadataResult(Dialect.GOOGLE_STANDARD_SQL)); futureParentHandlers = Logger.getLogger(AbstractFuture.class.getName()).getUseParentHandlers(); exceptionRunnableParentHandlers = diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/TransactionMockServerTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/TransactionMockServerTest.java index 45f68b11a5b6..bc27f4afbef1 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/TransactionMockServerTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/TransactionMockServerTest.java @@ -209,7 +209,7 @@ public void testBeginTransactionIsolationLevel() { SpannerPool.closeSpannerPool(); for (Dialect dialect : new Dialect[] {Dialect.POSTGRESQL, Dialect.GOOGLE_STANDARD_SQL}) { mockSpanner.putStatementResult( - MockSpannerServiceImpl.StatementResult.detectDialectResult(dialect)); + MockSpannerServiceImpl.StatementResult.detectMetadataResult(dialect)); try (Connection connection = super.createConnection()) { for (IsolationLevel isolationLevel : @@ -260,7 +260,7 @@ public void testBeginTransactionIsolationLevel() { public void testSetTransactionIsolationLevel() { SpannerPool.closeSpannerPool(); mockSpanner.putStatementResult( - MockSpannerServiceImpl.StatementResult.detectDialectResult(Dialect.POSTGRESQL)); + MockSpannerServiceImpl.StatementResult.detectMetadataResult(Dialect.POSTGRESQL)); try (Connection connection = super.createConnection()) { for (boolean autocommit : new boolean[] {true, false}) { diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpcTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpcTest.java index 275bbe66c20d..cf71e99b8df9 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpcTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpcTest.java @@ -829,6 +829,7 @@ public void testRouteToLeaderHeaderForReadWrite() { try (Spanner spanner = options.getService()) { final DatabaseClient databaseClient = spanner.getDatabaseClient(DatabaseId.of("[PROJECT]", "[INSTANCE]", "[DATABASE]")); + databaseClient.getDialect(); TransactionRunner runner = databaseClient.readWriteTransaction(); runner.run( transaction -> { diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GfeLatencyTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GfeLatencyTest.java index bcded26d685f..21b057c7a671 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GfeLatencyTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/GfeLatencyTest.java @@ -86,6 +86,7 @@ public class GfeLatencyTest { private static DatabaseClient databaseClientNoHeader; private static final String INSTANCE_ID = "fake-instance"; + private static final String INSTANCE_ID_NO_HEADER = "fake-instance-noheader"; private static final String DATABASE_ID = "fake-database"; private static final String PROJECT_ID = "fake-project"; @@ -181,7 +182,10 @@ public void sendHeaders(Metadata headers) { .start(); spannerNoHeader = createSpannerOptions(addressNoHeader, serverNoHeader).getService(); databaseClientNoHeader = - spannerNoHeader.getDatabaseClient(DatabaseId.of(PROJECT_ID, INSTANCE_ID, DATABASE_ID)); + spannerNoHeader.getDatabaseClient( + DatabaseId.of(PROJECT_ID, INSTANCE_ID_NO_HEADER, DATABASE_ID)); + databaseClient.getDialect(); + databaseClientNoHeader.getDialect(); } @AfterClass @@ -347,7 +351,7 @@ private long getMetric(View view, String method, boolean withOverride) { List tagValues = new java.util.ArrayList<>(); for (TagKey column : view.getColumns()) { if (column == SpannerRpcViews.INSTANCE_ID) { - tagValues.add(TagValue.create(INSTANCE_ID)); + tagValues.add(TagValue.create(withOverride ? INSTANCE_ID_NO_HEADER : INSTANCE_ID)); } else if (column == SpannerRpcViews.DATABASE_ID) { tagValues.add(TagValue.create(DATABASE_ID)); } else if (column == SpannerRpcViews.METHOD) {