diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerOptions.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerOptions.java index 4d5b9d903508..a7d77a7006f4 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerOptions.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerOptions.java @@ -332,6 +332,7 @@ static GcpChannelPoolOptions mergeWithDefaultChannelPoolOptions( private final Map mergedQueryOptions; private final CallCredentialsProvider callCredentialsProvider; + private final CallContextConfigurator callContextConfigurator; private final CloseableExecutorProvider asyncExecutorProvider; private final String compressorName; private final String emulatorHost; @@ -386,9 +387,11 @@ public interface CallCredentialsProvider { /** * {@link CallContextConfigurator} can be used to modify the {@link ApiCallContext} for one or - * more specific RPCs. This can be used to set specific timeout value for RPCs or use specific - * {@link CallCredentials} for an RPC. The {@link CallContextConfigurator} must be set as a value - * on the {@link Context} using the {@link SpannerOptions#CALL_CONTEXT_CONFIGURATOR_KEY} key. + * more specific RPCs. This can be used to set specific timeout values for RPCs or use specific + * {@link CallCredentials} for an RPC. The {@link CallContextConfigurator} can be configured at + * the client level using {@link Builder#setCallContextConfigurator(CallContextConfigurator)}, or + * on a per-call basis as a value on the {@link Context} using the {@link + * SpannerOptions#CALL_CONTEXT_CONFIGURATOR_KEY} key. * *

This API is meant for advanced users. Most users should instead use the {@link * SpannerCallContextTimeoutConfigurator} for setting timeouts per RPC. @@ -517,8 +520,10 @@ static SpannerMethod valueOf(ReqT request, MethodDescriptorExample usage: * @@ -574,43 +579,40 @@ public ApiCallContext configure( if (spannerMethod == null) { return null; } - switch (SpannerMethod.valueOf(request, method)) { + ApiCallContext callContext = context == null ? GrpcCallContext.createDefault() : context; + switch (spannerMethod) { case BATCH_UPDATE: return batchUpdateTimeout == null ? null - : GrpcCallContext.createDefault().withTimeoutDuration(batchUpdateTimeout); + : callContext.withTimeoutDuration(batchUpdateTimeout); case COMMIT: - return commitTimeout == null - ? null - : GrpcCallContext.createDefault().withTimeoutDuration(commitTimeout); + return commitTimeout == null ? null : callContext.withTimeoutDuration(commitTimeout); case EXECUTE_QUERY: return executeQueryTimeout == null ? null - : GrpcCallContext.createDefault() + : callContext .withTimeoutDuration(executeQueryTimeout) .withStreamWaitTimeoutDuration(executeQueryTimeout); case EXECUTE_UPDATE: return executeUpdateTimeout == null ? null - : GrpcCallContext.createDefault().withTimeoutDuration(executeUpdateTimeout); + : callContext.withTimeoutDuration(executeUpdateTimeout); case PARTITION_QUERY: return partitionQueryTimeout == null ? null - : GrpcCallContext.createDefault().withTimeoutDuration(partitionQueryTimeout); + : callContext.withTimeoutDuration(partitionQueryTimeout); case PARTITION_READ: return partitionReadTimeout == null ? null - : GrpcCallContext.createDefault().withTimeoutDuration(partitionReadTimeout); + : callContext.withTimeoutDuration(partitionReadTimeout); case READ: return readTimeout == null ? null - : GrpcCallContext.createDefault() + : callContext .withTimeoutDuration(readTimeout) .withStreamWaitTimeoutDuration(readTimeout); case ROLLBACK: - return rollbackTimeout == null - ? null - : GrpcCallContext.createDefault().withTimeoutDuration(rollbackTimeout); + return rollbackTimeout == null ? null : callContext.withTimeoutDuration(rollbackTimeout); default: } return null; @@ -1030,6 +1032,7 @@ protected SpannerOptions(Builder builder) { this.mergedQueryOptions = ImmutableMap.copyOf(merged); } callCredentialsProvider = builder.callCredentialsProvider; + callContextConfigurator = builder.callContextConfigurator; asyncExecutorProvider = builder.asyncExecutorProvider; compressorName = builder.compressorName; emulatorHost = builder.emulatorHost; @@ -1379,6 +1382,7 @@ private static Builder prepareBuilder(Builder builder) { private Duration grpcKeepAliveTime = Duration.ofSeconds(120); private Duration grpcKeepAliveTimeout = Duration.ofSeconds(20); private CallCredentialsProvider callCredentialsProvider; + private CallContextConfigurator callContextConfigurator; private CloseableExecutorProvider asyncExecutorProvider; private String compressorName; private String emulatorHost = System.getenv("SPANNER_EMULATOR_HOST"); @@ -1488,6 +1492,7 @@ protected Builder() { this.enableGrpcGcpOtelMetrics = options.enableGrpcGcpOtelMetrics; this.defaultQueryOptions = options.defaultQueryOptions; this.callCredentialsProvider = options.callCredentialsProvider; + this.callContextConfigurator = options.callContextConfigurator; this.grpcKeepAliveTime = options.grpcKeepAliveTime; this.grpcKeepAliveTimeout = options.grpcKeepAliveTimeout; this.asyncExecutorProvider = options.asyncExecutorProvider; @@ -1874,6 +1879,84 @@ public Builder setCallCredentialsProvider(CallCredentialsProvider callCredential return this; } + /** + * Configures a client-level {@link CallContextConfigurator} to apply custom gRPC options, + * timeouts, or credentials to RPCs executed by this Spanner client. + * + *

By default, Spanner clients allow customizing call options on individual requests using + * gRPC's thread-local {@link io.grpc.Context} with {@link #CALL_CONTEXT_CONFIGURATOR_KEY}. + * While useful for fine-grained per-RPC overrides, managing thread-local context can be + * cumbersome or error-prone in asynchronous, reactive, or multi-threaded pipelines where + * operations jump across threads. Setting a {@link CallContextConfigurator} here applies + * client-wide across all requests executed by this client instance without requiring + * thread-local context propagation. + * + *

This configurator applies to all RPCs executed by {@link DatabaseClient}, {@link Spanner}, + * {@link DatabaseAdminClient}, and {@link InstanceAdminClient} instances obtained from this + * client library. Note that raw GAPIC generated clients (such as {@link + * Spanner#createDatabaseAdminClient()} and {@link Spanner#createInstanceAdminClient()}) bypass + * this configurator and should be configured via {@link #setDatabaseAdminStubSettings} and + * {@link #setInstanceAdminStubSettings}. + * + *

Implementations of {@link CallContextConfigurator} configured at the client level must be + * thread-safe as they are shared across all concurrent operations executed by this client. + * + *

If both a client-level configurator and a thread-local configurator (via {@link + * #CALL_CONTEXT_CONFIGURATOR_KEY}) are present when an RPC is executed: + * + *

    + *
  1. The client-level configurator is evaluated first to establish the baseline call + * context. + *
  2. The thread-local configurator is evaluated next using that baseline context. + *
  3. Any options returned by the thread-local configurator are merged on top of the + * client-level options, allowing per-call configurations to override or extend + * client-level defaults. + *
+ * + *

Example: Configure a client-level stream wait timeout of 30 seconds for streaming SQL + * queries to detect stalled streams faster: + * + *

{@code
+     * SpannerOptions options =
+     *     SpannerOptions.newBuilder()
+     *         .setProjectId("my-project")
+     *         .setCallContextConfigurator(
+     *             new CallContextConfigurator() {
+     *               @Override
+     *               public  ApiCallContext configure(
+     *                   ApiCallContext context, ReqT request, MethodDescriptor method) {
+     *                 if (method == SpannerGrpc.getExecuteStreamingSqlMethod()) {
+     *                   return context.withStreamWaitTimeoutDuration(Duration.ofSeconds(30));
+     *                 }
+     *                 return null;
+     *               }
+     *             })
+     *         .build();
+     * }
+ * + *

You can also use {@link SpannerCallContextTimeoutConfigurator} if you only need to adjust + * standard timeouts across RPC types: + * + *

{@code
+     * SpannerOptions options =
+     *     SpannerOptions.newBuilder()
+     *         .setProjectId("my-project")
+     *         .setCallContextConfigurator(
+     *             SpannerCallContextTimeoutConfigurator.create()
+     *                 .withExecuteQueryTimeoutDuration(Duration.ofSeconds(30)))
+     *         .build();
+     * }
+ * + * @param callContextConfigurator the configurator to apply to all RPCs, or {@code null} to + * clear + * @return this {@link Builder} instance + */ + public Builder setCallContextConfigurator( + @Nullable CallContextConfigurator callContextConfigurator) { + this.callContextConfigurator = callContextConfigurator; + return this; + } + /** * Sets the compression to use for all gRPC calls. The compressor must be a valid name known in * the {@link CompressorRegistry}. This will enable compression both from the client to the @@ -2681,6 +2764,15 @@ public CallCredentialsProvider getCallCredentialsProvider() { return callCredentialsProvider; } + /** + * Returns the client-level {@link CallContextConfigurator} configured for this {@link + * SpannerOptions}, or {@code null} if none is set. + */ + @Nullable + public CallContextConfigurator getCallContextConfigurator() { + return callContextConfigurator; + } + private boolean usesNoCredentials() { // When JMH is enabled, we need to enable built-in metrics if (System.getProperty("jmh.enabled") != null diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpc.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpc.java index 7d969d39fc59..e76091148705 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpc.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/spi/v1/GapicSpannerRpc.java @@ -268,6 +268,9 @@ public class GapicSpannerRpc implements SpannerRpc { private static final String API_FILE = "grpc-gcp-apiconfig.json"; + private static final CallOptions.Key BASE_CONTEXT_MARKER_KEY = + CallOptions.Key.create("BASE_CONTEXT_MARKER_KEY"); + private final RequestIdCreator requestIdCreator = new RequestIdCreatorImpl(); private boolean rpcIsClosed; private final SpannerStub spannerStub; @@ -286,6 +289,7 @@ public class GapicSpannerRpc implements SpannerRpc { private final String projectName; private final SpannerMetadataProvider metadataProvider; private final CallCredentialsProvider callCredentialsProvider; + private final CallContextConfigurator callContextConfigurator; private final String compressorName; private final Duration waitTimeout = systemProperty(PROPERTY_TIMEOUT_SECONDS, DEFAULT_TIMEOUT_SECONDS); @@ -363,6 +367,7 @@ public GapicSpannerRpc(final SpannerOptions options) { headerProviderWithUserAgent.getHeaders(), internalHeaderProviderBuilder.getResourceHeaderKey()); this.callCredentialsProvider = options.getCallCredentialsProvider(); + this.callContextConfigurator = options.getCallContextConfigurator(); this.compressorName = options.getCompressorName(); this.leaderAwareRoutingEnabled = options.isLeaderAwareRoutingEnabled(); this.endToEndTracingEnabled = options.isEndToEndTracingEnabled(); @@ -2111,7 +2116,7 @@ public StreamingCall read( requestId, request.getSession(), request, - SpannerGrpc.getReadMethod(), + SpannerGrpc.getStreamingReadMethod(), routeToLeader); SpannerResponseObserver responseObserver = new SpannerResponseObserver(consumer); spannerStub.streamingReadCallable().call(request, responseObserver, context); @@ -2448,6 +2453,24 @@ GrpcCallContext newCallContext( MethodDescriptor method, boolean routeToLeader) { GrpcCallContext context = this.baseGrpcCallContext; + if (callCredentialsProvider != null) { + CallCredentials callCredentials = callCredentialsProvider.getCallCredentials(); + if (callCredentials != null) { + context = + context.withCallOptions(context.getCallOptions().withCallCredentials(callCredentials)); + } + } + + // 1. Sequentially evaluate client-level and thread-level configurators. + // The thread-level configurator receives the context after client-level modifications, + // allowing thread-scoped settings to override or extend client defaults. + context = applyConfigurators(context, request, method); + + // 2. Attach Spanner-internal routing options and headers to the final context. + // Doing this AFTER configurator evaluation guarantees that internal options (request ID, + // channel affinity) cannot be wiped out by configurators returning + // GrpcCallContext.createDefault(), + // and internal headers (resource prefix, route-to-leader) cannot be duplicated. Long affinity = options == null ? null : Option.CHANNEL_HINT.getLong(options); ChannelAffinityRef channelAffinityRef = options == null ? null : Option.CHANNEL_ID_AFFINITY.getChannelAffinityRef(options); @@ -2489,19 +2512,78 @@ GrpcCallContext newCallContext( if (routeToLeader && leaderAwareRoutingEnabled) { context = context.withExtraHeaders(metadataProvider.newRouteToLeaderHeader()); } - if (callCredentialsProvider != null) { - CallCredentials callCredentials = callCredentialsProvider.getCallCredentials(); - if (callCredentials != null) { - context = - context.withCallOptions(context.getCallOptions().withCallCredentials(callCredentials)); - } + if (compressorName != null && context.getCallOptions().getCompressor() == null) { + context = context.withCallOptions(context.getCallOptions().withCompression(compressorName)); } - CallContextConfigurator configurator = SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY.get(); - ApiCallContext apiCallContextFromContext = null; - if (configurator != null) { - apiCallContextFromContext = configurator.configure(context, request, method); + return context; + } + + private GrpcCallContext applyConfigurators( + GrpcCallContext context, ReqT request, MethodDescriptor method) { + if (method == null) { + return context; } - return (GrpcCallContext) context.merge(apiCallContextFromContext); + CallContextConfigurator threadConfigurator = SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY.get(); + if (this.callContextConfigurator == null && threadConfigurator == null) { + return context; + } + GrpcCallContext callContext = + context.withCallOptions( + context.getCallOptions().withOption(BASE_CONTEXT_MARKER_KEY, Boolean.TRUE)); + if (this.callContextConfigurator != null) { + callContext = + applySingleConfigurator(callContext, this.callContextConfigurator, request, method); + } + if (threadConfigurator != null) { + callContext = applySingleConfigurator(callContext, threadConfigurator, request, method); + } + return callContext; + } + + private static GrpcCallContext applySingleConfigurator( + GrpcCallContext base, + CallContextConfigurator configurator, + ReqT request, + MethodDescriptor method) { + ApiCallContext configured; + try { + configured = configurator.configure(base, request, method); + } catch (Throwable t) { + throw SpannerExceptionFactory.asSpannerException(t); + } + if (configured == null || configured == base) { + return base; + } + if (!(configured instanceof GrpcCallContext)) { + throw new IllegalArgumentException( + "context must be an instance of GrpcCallContext, but found " + + configured.getClass().getName()); + } + GrpcCallContext overlay = (GrpcCallContext) configured; + + // Check whether overlay was derived from base (retaining BASE_CONTEXT_MARKER_KEY) + // or is a standalone delta context (such as one created via GrpcCallContext.createDefault()). + boolean isDerived = + Boolean.TRUE.equals(overlay.getCallOptions().getOption(BASE_CONTEXT_MARKER_KEY)); + if (isDerived) { + return overlay; + } + + // Overlay is a standalone delta context. Merge it onto base. + GrpcCallContext merged = (GrpcCallContext) base.merge(overlay); + + // If the delta context did not set custom CallOptions, GAX's merge would replace base's + // CallOptions with CallOptions.DEFAULT. In that case, preserve base's CallOptions. + if (overlay.getCallOptions().equals(CallOptions.DEFAULT) + && !base.getCallOptions().equals(CallOptions.DEFAULT)) { + merged = merged.withCallOptions(base.getCallOptions()); + } else if (!Boolean.TRUE.equals(merged.getCallOptions().getOption(BASE_CONTEXT_MARKER_KEY))) { + merged = + merged.withCallOptions( + merged.getCallOptions().withOption(BASE_CONTEXT_MARKER_KEY, Boolean.TRUE)); + } + + return merged; } @Override diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpannerOptionsTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpannerOptionsTest.java index 25813260944f..4b754c74027f 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpannerOptionsTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SpannerOptionsTest.java @@ -42,6 +42,7 @@ import com.google.cloud.grpc.GcpChannelPrimer; import com.google.cloud.grpc.GcpManagedChannelOptions.GcpChannelPoolOptions; import com.google.cloud.spanner.SpannerOptions.Builder.DefaultReadWriteTransactionOptions; +import com.google.cloud.spanner.SpannerOptions.CallContextConfigurator; import com.google.cloud.spanner.SpannerOptions.FixedCloseableExecutorProvider; import com.google.cloud.spanner.SpannerOptions.SpannerCallContextTimeoutConfigurator; import com.google.cloud.spanner.admin.database.v1.stub.DatabaseAdminStubSettings; @@ -69,6 +70,7 @@ import com.google.spanner.v1.RollbackRequest; import com.google.spanner.v1.SpannerGrpc; import com.google.spanner.v1.TransactionOptions.IsolationLevel; +import io.grpc.MethodDescriptor; import io.opentelemetry.api.GlobalOpenTelemetry; import io.opentelemetry.api.OpenTelemetry; import io.opentelemetry.sdk.OpenTelemetrySdk; @@ -1649,4 +1651,37 @@ public void testGrpcKeepAliveTimeout() { IllegalArgumentException.class, () -> SpannerOptions.newBuilder().setGrpcKeepAliveTimeout(Duration.ofSeconds(-10))); } + + @Test + public void testCallContextConfigurator() { + SpannerOptions defaultOptions = + SpannerOptions.newBuilder() + .setProjectId("test-project") + .setCredentials(NoCredentials.getInstance()) + .build(); + assertNull(defaultOptions.getCallContextConfigurator()); + + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return null; + } + }; + SpannerOptions customOptions = + SpannerOptions.newBuilder() + .setProjectId("test-project") + .setCredentials(NoCredentials.getInstance()) + .setCallContextConfigurator(configurator) + .build(); + assertSame(configurator, customOptions.getCallContextConfigurator()); + + SpannerOptions optionsFromBuilder = customOptions.toBuilder().build(); + assertSame(configurator, optionsFromBuilder.getCallContextConfigurator()); + + SpannerOptions clearedOptions = + customOptions.toBuilder().setCallContextConfigurator(null).build(); + assertNull(clearedOptions.getCallContextConfigurator()); + } } 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 72c0f47243be..d1e07caa5108 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 @@ -16,6 +16,7 @@ package com.google.cloud.spanner.spi.v1; +import static com.google.cloud.spanner.XGoogSpannerRequestId.REQUEST_ID_CALL_OPTIONS_KEY; import static com.google.common.truth.Truth.assertThat; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertEquals; @@ -63,6 +64,7 @@ import com.google.cloud.spanner.SpannerExceptionFactory; import com.google.cloud.spanner.SpannerOptions; import com.google.cloud.spanner.SpannerOptions.CallContextConfigurator; +import com.google.cloud.spanner.SpannerOptions.SpannerCallContextTimeoutConfigurator; import com.google.cloud.spanner.SpannerOptionsHelper; import com.google.cloud.spanner.Statement; import com.google.cloud.spanner.TransactionRunner; @@ -73,15 +75,20 @@ import com.google.common.util.concurrent.Futures; import com.google.protobuf.ListValue; import com.google.rpc.ErrorInfo; +import com.google.spanner.admin.database.v1.DatabaseAdminGrpc; +import com.google.spanner.admin.database.v1.ListDatabaseRolesRequest; +import com.google.spanner.v1.CommitRequest; import com.google.spanner.v1.CreateSessionRequest; import com.google.spanner.v1.ExecuteSqlRequest; import com.google.spanner.v1.GetSessionRequest; +import com.google.spanner.v1.ReadRequest; import com.google.spanner.v1.ResultSetMetadata; import com.google.spanner.v1.Session; import com.google.spanner.v1.SpannerGrpc; import com.google.spanner.v1.StructType; import com.google.spanner.v1.StructType.Field; import com.google.spanner.v1.TypeCode; +import io.grpc.CallOptions; import io.grpc.Context; import io.grpc.Contexts; import io.grpc.ManagedChannelBuilder; @@ -109,6 +116,7 @@ import java.io.IOException; import java.lang.reflect.Array; import java.lang.reflect.Modifier; +import java.lang.reflect.Proxy; import java.net.InetSocketAddress; import java.net.URLEncoder; import java.time.Duration; @@ -685,6 +693,38 @@ public void testCallCredentialsProviderPreferenceAboveCredentials() { rpc.shutdown(); } + @Test + public void + testCallCredentialsProviderPreservedWhenConfiguratorReturnsDeltaContextWithCustomCallOptions() { + CallOptions.Key customKey = CallOptions.Key.create("customKey"); + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return GrpcCallContext.createDefault() + .withCallOptions(CallOptions.DEFAULT.withOption(customKey, "customVal")); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCredentials(STATIC_CREDENTIALS) + .setCallCredentialsProvider(() -> MoreCallCredentials.from(VARIABLE_CREDENTIALS)) + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + GetSessionRequest.getDefaultInstance(), + SpannerGrpc.getGetSessionMethod()); + assertNotNull(callContext.getCallOptions().getCredentials()); + assertEquals("customVal", callContext.getCallOptions().getOption(customKey)); + rpc.shutdown(); + } + @Test public void testCallCredentialsProviderReturnsNull() { SpannerOptions options = @@ -834,6 +874,122 @@ public ApiCallContext configure( }); } + @Test + public void testClientLevelCallContextConfiguratorEndToEnd() { + final TimeoutHolder timeoutHolder = new TimeoutHolder(); + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + if (request instanceof ExecuteSqlRequest + && method.equals(SpannerGrpc.getExecuteSqlMethod())) { + ExecuteSqlRequest sqlRequest = (ExecuteSqlRequest) request; + if (sqlRequest.getSeqno() > 0L) { + return context.withTimeoutDuration(timeoutHolder.timeout); + } + } + return null; + } + }; + + mockSpanner.setExecuteSqlExecutionTime(SimulatedExecutionTime.ofMinimumAndRandomTime(10, 0)); + SpannerOptions options = + createSpannerOptions().toBuilder().setCallContextConfigurator(configurator).build(); + try (Spanner customSpanner = options.getService()) { + DatabaseClient client = + customSpanner.getDatabaseClient(DatabaseId.of("[PROJECT]", "[INSTANCE]", "[DATABASE]")); + + // 1. A 1ns timeout causes a DEADLINE_EXCEEDED exception end-to-end. + timeoutHolder.timeout = Duration.ofNanos(1L); + SpannerException e = + assertThrows( + SpannerException.class, + () -> + client + .readWriteTransaction() + .run(transaction -> transaction.executeUpdate(UPDATE_FOO_STATEMENT))); + assertEquals(ErrorCode.DEADLINE_EXCEEDED, e.getErrorCode()); + + // 2. A longer timeout succeeds. + timeoutHolder.timeout = Duration.ofMinutes(1L); + long updateCount = + client + .readWriteTransaction() + .run(transaction -> transaction.executeUpdate(UPDATE_FOO_STATEMENT)); + assertEquals(1L, updateCount); + } + } + + @Test + public void testThreadLevelOverridesClientLevelCallContextConfiguratorEndToEnd() { + // Client-level configurator sets a 1-minute timeout which normally succeeds. + CallContextConfigurator clientConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + if (request instanceof ExecuteSqlRequest + && method.equals(SpannerGrpc.getExecuteSqlMethod())) { + ExecuteSqlRequest sqlRequest = (ExecuteSqlRequest) request; + if (sqlRequest.getSeqno() > 0L) { + return context.withTimeoutDuration(Duration.ofMinutes(1L)); + } + } + return null; + } + }; + + // Thread-level configurator overrides the timeout to 1ns. + CallContextConfigurator threadConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + if (request instanceof ExecuteSqlRequest + && method.equals(SpannerGrpc.getExecuteSqlMethod())) { + ExecuteSqlRequest sqlRequest = (ExecuteSqlRequest) request; + if (sqlRequest.getSeqno() > 0L) { + return context.withTimeoutDuration(Duration.ofNanos(1L)); + } + } + return null; + } + }; + + mockSpanner.setExecuteSqlExecutionTime(SimulatedExecutionTime.ofMinimumAndRandomTime(10, 0)); + SpannerOptions options = + createSpannerOptions().toBuilder().setCallContextConfigurator(clientConfigurator).build(); + try (Spanner customSpanner = options.getService()) { + DatabaseClient client = + customSpanner.getDatabaseClient(DatabaseId.of("[PROJECT]", "[INSTANCE]", "[DATABASE]")); + + // Without thread-level configurator, the client-level 1m timeout succeeds. + long updateCount = + client + .readWriteTransaction() + .run(transaction -> transaction.executeUpdate(UPDATE_FOO_STATEMENT)); + assertEquals(1L, updateCount); + + // With thread-level configurator, the 1ns timeout overrides client-level and causes + // DEADLINE_EXCEEDED. + Context context = + Context.current() + .withValue(SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY, threadConfigurator); + context.run( + () -> { + SpannerException e = + assertThrows( + SpannerException.class, + () -> + client + .readWriteTransaction() + .run(transaction -> transaction.executeUpdate(UPDATE_FOO_STATEMENT))); + assertEquals(ErrorCode.DEADLINE_EXCEEDED, e.getErrorCode()); + }); + } + } + @Test public void testNewCallContextWithNullRequestAndNullMethod() { SpannerOptions options = SpannerOptions.newBuilder().setProjectId("some-project").build(); @@ -842,6 +998,574 @@ public void testNewCallContextWithNullRequestAndNullMethod() { rpc.shutdown(); } + @Test + public void testNewCallContextWithClientLevelConfigurator() { + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + if (method.equals(SpannerGrpc.getExecuteStreamingSqlMethod())) { + return context + .withTimeoutDuration(Duration.ofSeconds(60)) + .withStreamWaitTimeoutDuration(Duration.ofSeconds(30)); + } + return null; + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(60), callContext.getTimeoutDuration()); + assertEquals(Duration.ofSeconds(30), callContext.getStreamWaitTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext.getExtraHeaders().get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithThreadLevelOverridesClientLevelConfigurator() { + CallContextConfigurator clientConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withTimeoutDuration(Duration.ofSeconds(60)); + } + }; + CallContextConfigurator threadConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withTimeoutDuration(Duration.ofSeconds(10)); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(clientConfigurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + Context context = + Context.current() + .withValue(SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY, threadConfigurator); + context.run( + () -> { + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(10), callContext.getTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext + .getExtraHeaders() + .get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + }); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithBothClientAndThreadLevelConfiguratorsMerged() { + CallContextConfigurator clientConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withStreamWaitTimeoutDuration(Duration.ofSeconds(30)); + } + }; + CallContextConfigurator threadConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withTimeoutDuration(Duration.ofSeconds(10)); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(clientConfigurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + Context context = + Context.current() + .withValue(SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY, threadConfigurator); + context.run( + () -> { + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(10), callContext.getTimeoutDuration()); + assertEquals(Duration.ofSeconds(30), callContext.getStreamWaitTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext + .getExtraHeaders() + .get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + }); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithSpannerCallContextTimeoutConfiguratorClientLevel() { + SpannerCallContextTimeoutConfigurator configurator = + SpannerCallContextTimeoutConfigurator.create() + .withExecuteQueryTimeoutDuration(Duration.ofSeconds(45)); + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(45), callContext.getTimeoutDuration()); + assertEquals(Duration.ofSeconds(45), callContext.getStreamWaitTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext.getExtraHeaders().get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithDerivedContextNoDuplicateHeaders() { + CallContextConfigurator clientConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withTimeoutDuration(Duration.ofSeconds(60)); + } + }; + CallContextConfigurator threadConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withTimeoutDuration(Duration.ofSeconds(10)); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .enableLeaderAwareRouting() + .setCallContextConfigurator(clientConfigurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + Context context = + Context.current() + .withValue(SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY, threadConfigurator); + context.run( + () -> { + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + /* requestId= */ null, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod(), + /* routeToLeader= */ true); + assertEquals(Duration.ofSeconds(10), callContext.getTimeoutDuration()); + assertEquals( + ImmutableList.of("true"), + callContext.getExtraHeaders().get("x-goog-spanner-route-to-leader")); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext + .getExtraHeaders() + .get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + }); + rpc.shutdown(); + } + + @Test + public void testNewCallContextPreservesCustomCallOptionsAndUserHeaders() { + CallOptions.Key customOptionKey = CallOptions.Key.create("customOptionKey"); + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return ((GrpcCallContext) context) + .withCallOptions( + ((GrpcCallContext) context) + .getCallOptions() + .withOption(customOptionKey, "customValue")) + .withExtraHeaders( + Collections.singletonMap("custom-header", ImmutableList.of("customValue"))); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals("customValue", callContext.getCallOptions().getOption(customOptionKey)); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + assertEquals( + ImmutableList.of("customValue"), callContext.getExtraHeaders().get("custom-header")); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext.getExtraHeaders().get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithInvalidContextThrowsIllegalArgumentException() { + ApiCallContext invalidContext = + (ApiCallContext) + Proxy.newProxyInstance( + ApiCallContext.class.getClassLoader(), + new Class[] {ApiCallContext.class}, + (proxy, method, args) -> null); + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return invalidContext; + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + assertThrows( + IllegalArgumentException.class, + () -> + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod())); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithConfiguratorReturningNullReturnsDefault() { + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return null; + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertNotNull(callContext); + assertNull(callContext.getTimeoutDuration()); + assertEquals(Duration.ofMinutes(30), callContext.getStreamWaitTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + assertEquals( + ImmutableList.of("projects/some-project"), + callContext.getExtraHeaders().get(ApiClientHeaderProvider.getDefaultResourceHeaderKey())); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithSpannerCallContextTimeoutConfiguratorForRead() { + SpannerCallContextTimeoutConfigurator configurator = + SpannerCallContextTimeoutConfigurator.create() + .withReadTimeoutDuration(Duration.ofSeconds(45)); + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ReadRequest.getDefaultInstance(), + SpannerGrpc.getStreamingReadMethod()); + assertEquals(Duration.ofSeconds(45), callContext.getTimeoutDuration()); + assertEquals(Duration.ofSeconds(45), callContext.getStreamWaitTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithSpannerCallContextTimeoutConfiguratorForCommit() { + SpannerCallContextTimeoutConfigurator configurator = + SpannerCallContextTimeoutConfigurator.create() + .withCommitTimeoutDuration(Duration.ofSeconds(20)); + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + CommitRequest.getDefaultInstance(), + SpannerGrpc.getCommitMethod()); + assertEquals(Duration.ofSeconds(20), callContext.getTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + rpc.shutdown(); + } + + @Test + public void + testNewCallContextPreservesCustomCallOptionsWhenThreadConfiguratorReturnsDeltaContext() { + CallOptions.Key customOptionKey = CallOptions.Key.create("customOptionKey"); + CallContextConfigurator clientConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return ((GrpcCallContext) context) + .withCallOptions( + ((GrpcCallContext) context) + .getCallOptions() + .withOption(customOptionKey, "customValue")); + } + }; + // Thread configurator returns a standalone delta context with default CallOptions. + CallContextConfigurator threadConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return GrpcCallContext.createDefault().withTimeoutDuration(Duration.ofSeconds(10)); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(clientConfigurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + Context context = + Context.current() + .withValue(SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY, threadConfigurator); + context.run( + () -> { + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(10), callContext.getTimeoutDuration()); + assertEquals("customValue", callContext.getCallOptions().getOption(customOptionKey)); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + }); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithDeltaContextCustomIdleTimeout() { + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return GrpcCallContext.createDefault() + .withStreamIdleTimeoutDuration(Duration.ofSeconds(15)); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(15), callContext.getStreamIdleTimeoutDuration()); + assertEquals(Duration.ofMinutes(30), callContext.getStreamWaitTimeoutDuration()); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + rpc.shutdown(); + } + + @Test + public void testNewCallContextClientDeltaWithCustomCallOptionsAndThreadDerivedContext() { + CallOptions.Key clientKey = CallOptions.Key.create("clientKey"); + CallContextConfigurator clientConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return GrpcCallContext.createDefault() + .withCallOptions(CallOptions.DEFAULT.withOption(clientKey, "clientValue")) + .withExtraHeaders( + Collections.singletonMap("custom-header", ImmutableList.of("headerValue"))); + } + }; + CallContextConfigurator threadConfigurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + return context.withTimeoutDuration(Duration.ofSeconds(15)); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(clientConfigurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + Context context = + Context.current() + .withValue(SpannerOptions.CALL_CONTEXT_CONFIGURATOR_KEY, threadConfigurator); + context.run( + () -> { + GrpcCallContext callContext = + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod()); + assertEquals(Duration.ofSeconds(15), callContext.getTimeoutDuration()); + assertEquals("clientValue", callContext.getCallOptions().getOption(clientKey)); + assertEquals( + ImmutableList.of("headerValue"), callContext.getExtraHeaders().get("custom-header")); + assertNotNull(callContext.getCallOptions().getOption(REQUEST_ID_CALL_OPTIONS_KEY)); + }); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithNullMethodDoesNotThrowNpe() { + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + if (method.equals(SpannerGrpc.getExecuteStreamingSqlMethod())) { + return context.withTimeoutDuration(Duration.ofSeconds(10)); + } + return null; + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + GrpcCallContext callContext = rpc.newCallContext(optionsMap, "/some/resource", null, null); + assertNotNull(callContext); + rpc.shutdown(); + } + + @Test + public void testNewCallContextConfiguratorThrowsExceptionWrapsInSpannerException() { + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + throw new RuntimeException("custom configurator failure"); + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + SpannerException e = + assertThrows( + SpannerException.class, + () -> + rpc.newCallContext( + optionsMap, + "/some/resource", + ExecuteSqlRequest.getDefaultInstance(), + SpannerGrpc.getExecuteStreamingSqlMethod())); + assertTrue(e.getMessage().contains("custom configurator failure")); + rpc.shutdown(); + } + + @Test + public void testNewCallContextWithAdminMethodAppliesConfigurator() { + CallContextConfigurator configurator = + new CallContextConfigurator() { + @Override + public ApiCallContext configure( + ApiCallContext context, ReqT request, MethodDescriptor method) { + if (method == DatabaseAdminGrpc.getListDatabaseRolesMethod()) { + return context.withTimeoutDuration(Duration.ofSeconds(42)); + } + return null; + } + }; + SpannerOptions options = + SpannerOptions.newBuilder() + .setProjectId("some-project") + .setCallContextConfigurator(configurator) + .build(); + GapicSpannerRpc rpc = new GapicSpannerRpc(options, false); + ListDatabaseRolesRequest request = + ListDatabaseRolesRequest.newBuilder() + .setParent("projects/p/instances/i/databases/d") + .build(); + GrpcCallContext callContext = + rpc.newCallContext( + /* options= */ null, + "projects/p/instances/i/databases/d", + request, + DatabaseAdminGrpc.getListDatabaseRolesMethod()); + assertEquals(Duration.ofSeconds(42), callContext.getTimeoutDuration()); + rpc.shutdown(); + } + @Test public void testNewCallContextWithGrpcGcpUsesChannelAffinityRefWithoutDcp() { SpannerOptions options =