Skip to content

Commit 2b41e1c

Browse files
committed
Fix: suggested changes
1 parent 99a90c4 commit 2b41e1c

10 files changed

Lines changed: 161 additions & 27 deletions

File tree

api/src/main/java/io/grpc/ManagedChannelBuilder.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -666,7 +666,7 @@ public <X> T setNameResolverArg(NameResolver.Args.Key<X> key, X value) {
666666
* Sets a configurator that will be applied to all internal child channels created by this
667667
* channel.
668668
*
669-
* <p>This allows injecting configuration (like credentials, interceptors, or flow control)
669+
* <p>This allows injecting universal configuration (like interceptors)
670670
* into auxiliary channels created by gRPC infrastructure, such as xDS control plane connections.
671671
*
672672
* @param channelConfigurator the configurator to apply.

api/src/main/java/io/grpc/NameResolver.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -478,7 +478,7 @@ public ChannelLogger getChannelLogger() {
478478
*
479479
* @since 1.83.0
480480
*/
481-
@Internal
481+
@ExperimentalApi("https://github.com/grpc/grpc-java/issues/12574")
482482
public ChannelConfigurator getChildChannelConfigurator() {
483483
return channelConfigurator;
484484
}
@@ -561,6 +561,9 @@ public Builder toBuilder() {
561561
builder.setOverrideAuthority(overrideAuthority);
562562
builder.setMetricRecorder(metricRecorder);
563563
builder.setNameResolverRegistry(nameResolverRegistry);
564+
if (channelConfigurator != null) {
565+
builder.setChildChannelConfigurator(channelConfigurator);
566+
}
564567
builder.customArgs = cloneCustomArgs(customArgs);
565568
return builder;
566569
}
@@ -712,6 +715,7 @@ public Builder setNameResolverRegistry(NameResolverRegistry registry) {
712715
*
713716
* @since 1.83.0
714717
*/
718+
@ExperimentalApi("https://github.com/grpc/grpc-java/issues/12574")
715719
public Builder setChildChannelConfigurator(ChannelConfigurator channelConfigurator) {
716720
this.channelConfigurator = checkNotNull(channelConfigurator, "channelConfigurator");
717721
return this;

api/src/test/java/io/grpc/NameResolverTest.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import static org.mockito.Mockito.verify;
2323

2424
import com.google.common.base.Objects;
25+
import io.grpc.ChannelConfigurator;
2526
import io.grpc.NameResolver.ConfigOrError;
2627
import io.grpc.NameResolver.Listener2;
2728
import io.grpc.NameResolver.ResolutionResult;
@@ -72,6 +73,8 @@ public class NameResolverTest {
7273
private final int customArgValue = 42;
7374
@Mock NameResolver.Listener mockListener;
7475

76+
private final ChannelConfigurator channelConfigurator = builder -> { };
77+
7578
@Test
7679
public void args() {
7780
NameResolver.Args args = createArgs();
@@ -84,6 +87,7 @@ public void args() {
8487
assertThat(args.getOffloadExecutor()).isSameInstanceAs(executor);
8588
assertThat(args.getOverrideAuthority()).isSameInstanceAs(overrideAuthority);
8689
assertThat(args.getMetricRecorder()).isSameInstanceAs(metricRecorder);
90+
assertThat(args.getChildChannelConfigurator()).isSameInstanceAs(channelConfigurator);
8791
assertThat(args.getArg(FOO_ARG_KEY)).isEqualTo(customArgValue);
8892
assertThat(args.getArg(BAR_ARG_KEY)).isNull();
8993

@@ -97,6 +101,7 @@ public void args() {
97101
assertThat(args2.getOffloadExecutor()).isSameInstanceAs(executor);
98102
assertThat(args2.getOverrideAuthority()).isSameInstanceAs(overrideAuthority);
99103
assertThat(args.getMetricRecorder()).isSameInstanceAs(metricRecorder);
104+
assertThat(args2.getChildChannelConfigurator()).isSameInstanceAs(channelConfigurator);
100105
assertThat(args.getArg(FOO_ARG_KEY)).isEqualTo(customArgValue);
101106
assertThat(args.getArg(BAR_ARG_KEY)).isNull();
102107

@@ -105,7 +110,6 @@ public void args() {
105110
}
106111

107112
private NameResolver.Args createArgs() {
108-
ChannelConfigurator channelConfigurator = builder -> { };
109113
return NameResolver.Args.newBuilder()
110114
.setDefaultPort(defaultPort)
111115
.setProxyDetector(proxyDetector)
@@ -116,8 +120,8 @@ private NameResolver.Args createArgs() {
116120
.setOffloadExecutor(executor)
117121
.setOverrideAuthority(overrideAuthority)
118122
.setMetricRecorder(metricRecorder)
119-
.setArg(FOO_ARG_KEY, customArgValue)
120123
.setChildChannelConfigurator(channelConfigurator)
124+
.setArg(FOO_ARG_KEY, customArgValue)
121125
.build();
122126
}
123127

core/src/test/java/io/grpc/internal/ManagedChannelImplTest.java

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@
6868
import io.grpc.CallCredentials.RequestInfo;
6969
import io.grpc.CallOptions;
7070
import io.grpc.Channel;
71+
import io.grpc.ChannelConfigurator;
7172
import io.grpc.ChannelCredentials;
7273
import io.grpc.ChannelLogger;
7374
import io.grpc.ClientCall;
@@ -496,6 +497,40 @@ public void immediateDeadlineExceeded() {
496497
assertSame(Status.DEADLINE_EXCEEDED.getCode(), status.getCode());
497498
}
498499

500+
@Test
501+
public void childChannelConfigurator_passedToNameResolverArgs() {
502+
ChannelConfigurator configurator = builder -> { };
503+
channelBuilder.childChannelConfigurator(configurator);
504+
AtomicReference<NameResolver.Args> actualArgs = new AtomicReference<>();
505+
channelBuilder.nameResolverRegistry.register(new NameResolverProvider() {
506+
@Override
507+
public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
508+
actualArgs.set(args);
509+
NameResolver resolver = mock(NameResolver.class);
510+
when(resolver.getServiceAuthority()).thenReturn("test.example.com");
511+
return resolver;
512+
}
513+
514+
@Override
515+
public String getDefaultScheme() {
516+
return expectedUri.getScheme();
517+
}
518+
519+
@Override
520+
protected boolean isAvailable() {
521+
return true;
522+
}
523+
524+
@Override
525+
protected int priority() {
526+
return 10;
527+
}
528+
});
529+
createChannel();
530+
assertNotNull(actualArgs.get());
531+
assertSame(configurator, actualArgs.get().getChildChannelConfigurator());
532+
}
533+
499534
@Test
500535
public void startCallBeforeNameResolution() throws Exception {
501536
FakeNameResolverFactory nameResolverFactory =

googleapis/src/main/java/io/grpc/googleapis/GoogleCloudToProdNameResolver.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import com.google.common.collect.ImmutableMap;
2525
import com.google.common.io.CharStreams;
2626
import com.google.errorprone.annotations.concurrent.GuardedBy;
27+
import io.grpc.ChannelConfigurator;
2728
import io.grpc.MetricRecorder;
2829
import io.grpc.NameResolver;
2930
import io.grpc.NameResolverRegistry;
@@ -108,6 +109,7 @@ private static synchronized BootstrapInfo getBootstrapInfo(boolean isForcedXds)
108109
private final Resource<Executor> executorResource;
109110
private final String target;
110111
private final MetricRecorder metricRecorder;
112+
private final ChannelConfigurator channelConfigurator;
111113
private final NameResolver delegate;
112114
private final boolean usingExecutorResource;
113115
private final boolean forceXds;
@@ -160,6 +162,7 @@ private static synchronized BootstrapInfo getBootstrapInfo(boolean isForcedXds)
160162
}
161163
target = targetUri.toString();
162164
metricRecorder = args.getMetricRecorder();
165+
channelConfigurator = args.getChildChannelConfigurator();
163166
delegate = checkNotNull(nameResolverFactory, "nameResolverFactory").newNameResolver(
164167
targetUri, args);
165168
executor = args.getOffloadExecutor();
@@ -211,6 +214,7 @@ private static synchronized BootstrapInfo getBootstrapInfo(boolean isForcedXds)
211214
targetUri = modifiedTargetBuilder.build();
212215
target = targetUri.toString();
213216
metricRecorder = args.getMetricRecorder();
217+
channelConfigurator = args.getChildChannelConfigurator();
214218
delegate =
215219
checkNotNull(nameResolverFactory, "nameResolverFactory").newNameResolver(targetUri, args);
216220
executor = args.getOffloadExecutor();
@@ -278,7 +282,7 @@ public void run() {
278282
public void run() {
279283
if (!shutdown && finalBootstrapInfo != null) {
280284
xdsClientPool = InternalSharedXdsClientPoolProvider.getOrCreate(
281-
target, finalBootstrapInfo, metricRecorder, null);
285+
target, finalBootstrapInfo, metricRecorder, null, channelConfigurator);
282286
xdsClient = xdsClientPool.getObject();
283287
delegate.start(listener);
284288
succeeded = true;

googleapis/src/test/java/io/grpc/googleapis/GoogleCloudToProdNameResolverTest.java

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import static org.mockito.Mockito.when;
2323

2424
import com.google.common.collect.Iterables;
25+
import io.grpc.ChannelConfigurator;
2526
import io.grpc.ChannelLogger;
2627
import io.grpc.MetricRecorder;
2728
import io.grpc.NameResolver;
@@ -102,6 +103,7 @@ public void close(Executor instance) {}
102103
private final Map<String, NameResolver> delegatedResolver = new HashMap<>();
103104
private final Map<String, URI> delegatedUri = new HashMap<>();
104105
private final Map<String, Uri> delegatedRfcUri = new HashMap<>();
106+
private final Map<String, Args> delegatedArgs = new HashMap<>();
105107

106108
@Mock
107109
private NameResolver.Listener2 mockListener;
@@ -236,6 +238,22 @@ public void notOnGcpButForceXds_KeyValueTrue_DelegateToXds() {
236238
}
237239

238240

241+
@Test
242+
public void childChannelConfigurator_passedToDelegatedResolver() {
243+
GoogleCloudToProdNameResolver.isOnGcp = false;
244+
ChannelConfigurator configurator = builder -> { };
245+
Args customArgs = args.toBuilder().setChildChannelConfigurator(configurator).build();
246+
resolver = enableRfc3986UrisParam
247+
? new GoogleCloudToProdNameResolver(
248+
Uri.create(TARGET_URI), customArgs, fakeExecutorResource, nsRegistry.asFactory())
249+
: new GoogleCloudToProdNameResolver(
250+
URI.create(TARGET_URI), customArgs, fakeExecutorResource, nsRegistry.asFactory());
251+
resolver.start(mockListener);
252+
assertThat(delegatedArgs.keySet()).containsExactly("dns");
253+
assertThat(delegatedArgs.get("dns").getChildChannelConfigurator())
254+
.isSameInstanceAs(configurator);
255+
}
256+
239257
@Test
240258
public void notOnGcpButForceXds_WithMultipleParams_DelegateToXds() {
241259
GoogleCloudToProdNameResolver.isOnGcp = false;
@@ -338,6 +356,7 @@ private FakeNsProvider(String scheme) {
338356
public NameResolver newNameResolver(URI targetUri, Args args) {
339357
if (scheme.equals(targetUri.getScheme())) {
340358
delegatedUri.put(scheme, targetUri);
359+
delegatedArgs.put(scheme, args);
341360
NameResolver resolver = mock(NameResolver.class);
342361
delegatedResolver.put(scheme, resolver);
343362
return resolver;
@@ -349,6 +368,7 @@ public NameResolver newNameResolver(URI targetUri, Args args) {
349368
public NameResolver newNameResolver(Uri targetUri, Args args) {
350369
if (scheme.equals(targetUri.getScheme())) {
351370
delegatedRfcUri.put(scheme, targetUri);
371+
delegatedArgs.put(scheme, args);
352372
NameResolver resolver = mock(NameResolver.class);
353373
delegatedResolver.put(scheme, resolver);
354374
return resolver;

xds/src/main/java/io/grpc/xds/InternalSharedXdsClientPoolProvider.java

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package io.grpc.xds;
1818

1919
import io.grpc.CallCredentials;
20+
import io.grpc.ChannelConfigurator;
2021
import io.grpc.Internal;
2122
import io.grpc.MetricRecorder;
2223
import io.grpc.internal.ObjectPool;
@@ -85,9 +86,15 @@ public static ObjectPool<XdsClient> getOrCreate(
8586
public static XdsClientResult getOrCreate(
8687
String target, BootstrapInfo bootstrapInfo, MetricRecorder metricRecorder,
8788
CallCredentials transportCallCredentials) {
89+
return getOrCreate(target, bootstrapInfo, metricRecorder, transportCallCredentials, null);
90+
}
91+
92+
public static XdsClientResult getOrCreate(
93+
String target, BootstrapInfo bootstrapInfo, MetricRecorder metricRecorder,
94+
CallCredentials transportCallCredentials, ChannelConfigurator channelConfigurator) {
8895
return new XdsClientResult(SharedXdsClientPoolProvider.getDefaultProvider()
8996
.getOrCreate(target, bootstrapInfo, metricRecorder, transportCallCredentials,
90-
null));
97+
channelConfigurator));
9198
}
9299

93100
/**

xds/src/test/java/io/grpc/xds/FakeControlPlaneXdsOtelIntegrationTest.java

Lines changed: 41 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import io.grpc.ManagedChannel;
2727
import io.grpc.ManagedChannelBuilder;
2828
import io.grpc.opentelemetry.GrpcOpenTelemetry;
29+
import io.grpc.testing.GrpcCleanupRule;
2930
import io.grpc.testing.protobuf.SimpleRequest;
3031
import io.grpc.testing.protobuf.SimpleServiceGrpc;
3132
import io.opentelemetry.api.OpenTelemetry;
@@ -44,15 +45,19 @@
4445

4546
/**
4647
* xDS + OpenTelemetry E2E integration test using a fake control plane.
48+
* This class is skipped from Bazel builds because Bazel doesn't compile the
49+
* grpc-opentelemetry module.
4750
*/
4851
@RunWith(Parameterized.class)
4952
public class FakeControlPlaneXdsOtelIntegrationTest {
5053

5154
@Rule(order = 0)
52-
public ControlPlaneRule controlPlane = new ControlPlaneRule();
55+
public final GrpcCleanupRule grpcCleanupRule = new GrpcCleanupRule();
5356
@Rule(order = 1)
54-
public DataPlaneRule dataPlane = new DataPlaneRule(controlPlane);
57+
public final ControlPlaneRule controlPlane = new ControlPlaneRule();
5558
@Rule(order = 2)
59+
public final DataPlaneRule dataPlane = new DataPlaneRule(controlPlane);
60+
@Rule(order = 3)
5661
public final FlagResetRule flagResetRule = new FlagResetRule();
5762

5863
@Parameters(name = "enableRfc3986UrisParam={0}")
@@ -62,6 +67,21 @@ public static Iterable<Object[]> data() {
6267

6368
@Parameter public boolean enableRfc3986UrisParam;
6469

70+
@Test
71+
public void testInMemoryMetricReader() {
72+
InMemoryMetricReader metricReader = InMemoryMetricReader.create();
73+
SdkMeterProvider meterProvider = SdkMeterProvider.builder()
74+
.registerMetricReader(metricReader)
75+
.build();
76+
io.opentelemetry.api.metrics.LongCounter counter = meterProvider
77+
.meterBuilder("test-scope")
78+
.build()
79+
.counterBuilder("test-counter")
80+
.build();
81+
counter.add(10);
82+
assertThat(metricReader.collectAllMetrics()).isNotEmpty();
83+
}
84+
6585
@Before
6686
public void setupRfc3986UrisFeatureFlag() throws Exception {
6787
flagResetRule.setFlagForTest(
@@ -88,32 +108,32 @@ public void configureChannelBuilder(ManagedChannelBuilder<?> builder) {
88108
}
89109
};
90110

91-
ManagedChannel channel = Grpc.newChannelBuilder("test-xds:///test-server",
111+
ManagedChannel channel = grpcCleanupRule.register(
112+
Grpc.newChannelBuilder("test-xds:///test-server",
92113
InsecureChannelCredentials.create())
93114
.childChannelConfigurator(configurator)
94-
.build();
115+
.build());
95116

96-
try {
97-
SimpleServiceGrpc.SimpleServiceBlockingStub blockingStub = SimpleServiceGrpc.newBlockingStub(
98-
channel);
99-
blockingStub.unaryRpc(SimpleRequest.getDefaultInstance());
117+
SimpleServiceGrpc.SimpleServiceBlockingStub blockingStub =
118+
SimpleServiceGrpc.newBlockingStub(channel);
119+
blockingStub.unaryRpc(SimpleRequest.getDefaultInstance());
100120

101-
boolean hasMetrics = false;
102-
for (int i = 0; i < 20; i++) {
103-
for (MetricData metric : metricReader.collectAllMetrics()) {
104-
if (metric.getName().startsWith("grpc.client.")) {
105-
hasMetrics = true;
106-
break;
107-
}
108-
}
109-
if (hasMetrics) {
121+
// Verify that OpenTelemetry metrics specifically from the xDS Control Plane ADS stream
122+
// successfully propagated, method name is recorded as 'other' because dynamic descriptors
123+
// are not sampled to local tracing by default.
124+
boolean foundXdsMetrics = false;
125+
for (int i = 0; i < 50; i++) {
126+
for (MetricData metric : metricReader.collectAllMetrics()) {
127+
if (metric.toString().contains("grpc.client.") && metric.toString().contains("other")) {
128+
foundXdsMetrics = true;
110129
break;
111130
}
112-
Thread.sleep(100);
113131
}
114-
assertThat(hasMetrics).isTrue();
115-
} finally {
116-
channel.shutdownNow();
132+
if (foundXdsMetrics) {
133+
break;
134+
}
135+
Thread.sleep(100);
117136
}
137+
assertThat(foundXdsMetrics).isTrue();
118138
}
119139
}

xds/src/test/java/io/grpc/xds/XdsServerBuilderTest.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -349,4 +349,15 @@ public void start_passesChannelConfiguratorToClientPoolFactory() throws Exceptio
349349
verify(mockPoolFactory).getOrCreate(
350350
any(), any(), any(), eq(configurer));
351351
}
352+
353+
@Test
354+
public void childChannelConfigurator_nullThrows() throws IOException {
355+
buildBuilder(null);
356+
try {
357+
builder.childChannelConfigurator(null);
358+
fail("exception expected");
359+
} catch (NullPointerException expected) {
360+
assertThat(expected).hasMessageThat().contains("channelConfigurator");
361+
}
362+
}
352363
}

0 commit comments

Comments
 (0)