Skip to content

Commit cbb9667

Browse files
committed
Supports injecting custom ResourceNameResolvers
1 parent da2213c commit cbb9667

3 files changed

Lines changed: 210 additions & 16 deletions

File tree

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

Lines changed: 41 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -37,10 +37,14 @@
3737
import io.grpc.netty.InternalProtocolNegotiator;
3838
import io.grpc.netty.NettyServerBuilder;
3939
import io.grpc.xds.FilterChainMatchingProtocolNegotiators.FilterChainMatchingNegotiatorServerFactory;
40+
import java.net.InetSocketAddress;
41+
import java.net.SocketAddress;
4042
import java.util.Map;
4143
import java.util.concurrent.TimeUnit;
4244
import java.util.concurrent.atomic.AtomicBoolean;
45+
import java.util.function.Function;
4346
import java.util.logging.Logger;
47+
import javax.annotation.Nullable;
4448

4549
/**
4650
* A version of {@link ServerBuilder} to create xDS managed servers.
@@ -57,6 +61,7 @@ public final class XdsServerBuilder extends ForwardingServerBuilder<XdsServerBui
5761
private XdsClientPoolFactory xdsClientPoolFactory =
5862
SharedXdsClientPoolProvider.getDefaultProvider();
5963
private Map<String, ?> bootstrapOverride;
64+
@Nullable private Function<String, String> ldsResourceNameResolver;
6065
private long drainGraceTime = 10;
6166
private TimeUnit drainGraceTimeUnit = TimeUnit.MINUTES;
6267
private ChannelConfigurator channelConfigurator = builder -> { };
@@ -134,6 +139,23 @@ public static XdsServerBuilder forPort(int port, ServerCredentials serverCredent
134139
return new XdsServerBuilder(nettyDelegate, port);
135140
}
136141

142+
/** Creates a gRPC server builder for the given address. */
143+
public static XdsServerBuilder forAddress(
144+
SocketAddress address, ServerCredentials serverCredentials) {
145+
checkNotNull(serverCredentials, "serverCredentials");
146+
InternalProtocolNegotiator.ServerFactory originalNegotiatorFactory =
147+
InternalNettyServerCredentials.toNegotiator(serverCredentials);
148+
ServerCredentials wrappedCredentials =
149+
InternalNettyServerCredentials.create(
150+
new FilterChainMatchingNegotiatorServerFactory(originalNegotiatorFactory));
151+
NettyServerBuilder nettyDelegate = NettyServerBuilder.forAddress(address, wrappedCredentials);
152+
int port = 0;
153+
if (address instanceof InetSocketAddress inetSocketAddress) {
154+
port = inetSocketAddress.getPort();
155+
}
156+
return new XdsServerBuilder(nettyDelegate, port);
157+
}
158+
137159
@Override
138160
public Server build() {
139161
checkState(isServerBuilt.compareAndSet(false, true), "Server already built!");
@@ -144,11 +166,28 @@ public Server build() {
144166
builder.set(ATTR_DRAIN_GRACE_NANOS, drainGraceTimeUnit.toNanos(drainGraceTime));
145167
}
146168
InternalNettyServerBuilder.eagAttributes(delegate, builder.build());
147-
return new XdsServerWrapper("0.0.0.0:" + port, delegate, xdsServingStatusListener,
148-
filterChainSelectorManager, xdsClientPoolFactory, bootstrapOverride, filterRegistry,
169+
return new XdsServerWrapper(
170+
"0.0.0.0:" + port,
171+
delegate,
172+
xdsServingStatusListener,
173+
filterChainSelectorManager,
174+
xdsClientPoolFactory,
175+
bootstrapOverride,
176+
ldsResourceNameResolver,
177+
filterRegistry,
149178
this.channelConfigurator);
150179
}
151180

181+
/**
182+
* Provides a function that takes the listening address and returns the LDS resource name. When
183+
* provided, this overrides the server_listener_resource_name_template in the bootstrap.
184+
*/
185+
public XdsServerBuilder ldsResourceNameResolver(
186+
Function<String, String> ldsResourceNameResolver) {
187+
this.ldsResourceNameResolver = checkNotNull(ldsResourceNameResolver, "ldsResourceNameResolver");
188+
return this;
189+
}
190+
152191
@VisibleForTesting
153192
XdsServerBuilder xdsClientPoolFactory(XdsClientPoolFactory xdsClientPoolFactory) {
154193
this.xdsClientPoolFactory = checkNotNull(xdsClientPoolFactory, "xdsClientPoolFactory");

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

Lines changed: 125 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,7 @@
7878
import java.util.concurrent.TimeUnit;
7979
import java.util.concurrent.atomic.AtomicBoolean;
8080
import java.util.concurrent.atomic.AtomicReference;
81+
import java.util.function.Function;
8182
import java.util.logging.Level;
8283
import java.util.logging.Logger;
8384
import javax.annotation.Nullable;
@@ -108,6 +109,7 @@ public void uncaughtException(Thread t, Throwable e) {
108109
private final ThreadSafeRandom random = ThreadSafeRandomImpl.instance;
109110
private final XdsClientPoolFactory xdsClientPoolFactory;
110111
private final @Nullable Map<String, ?> bootstrapOverride;
112+
private final @Nullable Function<String, String> ldsResourceNameResolver;
111113
private final XdsServingStatusListener listener;
112114
private final FilterChainSelectorManager filterChainSelectorManager;
113115
private final AtomicBoolean started = new AtomicBoolean(false);
@@ -139,6 +141,7 @@ public void uncaughtException(Thread t, Throwable e) {
139141
FilterChainSelectorManager filterChainSelectorManager,
140142
XdsClientPoolFactory xdsClientPoolFactory,
141143
@Nullable Map<String, ?> bootstrapOverride,
144+
@Nullable Function<String, String> ldsResourceNameResolver,
142145
FilterRegistry filterRegistry,
143146
ChannelConfigurator channelConfigurator) {
144147
this(
@@ -148,6 +151,30 @@ public void uncaughtException(Thread t, Throwable e) {
148151
filterChainSelectorManager,
149152
xdsClientPoolFactory,
150153
bootstrapOverride,
154+
ldsResourceNameResolver,
155+
filterRegistry,
156+
SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE),
157+
channelConfigurator);
158+
sharedTimeService = true;
159+
}
160+
161+
XdsServerWrapper(
162+
String listenerAddress,
163+
ServerBuilder<?> delegateBuilder,
164+
XdsServingStatusListener listener,
165+
FilterChainSelectorManager filterChainSelectorManager,
166+
XdsClientPoolFactory xdsClientPoolFactory,
167+
@Nullable Map<String, ?> bootstrapOverride,
168+
FilterRegistry filterRegistry,
169+
ChannelConfigurator channelConfigurator) {
170+
this(
171+
listenerAddress,
172+
delegateBuilder,
173+
listener,
174+
filterChainSelectorManager,
175+
xdsClientPoolFactory,
176+
bootstrapOverride,
177+
null,
151178
filterRegistry,
152179
SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE),
153180
channelConfigurator);
@@ -169,8 +196,34 @@ public void uncaughtException(Thread t, Throwable e) {
169196
filterChainSelectorManager,
170197
xdsClientPoolFactory,
171198
bootstrapOverride,
199+
null,
172200
filterRegistry,
201+
SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE),
173202
builder -> { });
203+
sharedTimeService = true;
204+
}
205+
206+
XdsServerWrapper(
207+
String listenerAddress,
208+
ServerBuilder<?> delegateBuilder,
209+
XdsServingStatusListener listener,
210+
FilterChainSelectorManager filterChainSelectorManager,
211+
XdsClientPoolFactory xdsClientPoolFactory,
212+
@Nullable Map<String, ?> bootstrapOverride,
213+
@Nullable Function<String, String> ldsResourceNameResolver,
214+
FilterRegistry filterRegistry) {
215+
this(
216+
listenerAddress,
217+
delegateBuilder,
218+
listener,
219+
filterChainSelectorManager,
220+
xdsClientPoolFactory,
221+
bootstrapOverride,
222+
ldsResourceNameResolver,
223+
filterRegistry,
224+
SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE),
225+
builder -> {});
226+
sharedTimeService = true;
174227
}
175228

176229
@VisibleForTesting
@@ -190,11 +243,36 @@ public void uncaughtException(Thread t, Throwable e) {
190243
filterChainSelectorManager,
191244
xdsClientPoolFactory,
192245
bootstrapOverride,
246+
null,
193247
filterRegistry,
194248
timeService,
195249
builder -> { });
196250
}
197251

252+
@VisibleForTesting
253+
XdsServerWrapper(
254+
String listenerAddress,
255+
ServerBuilder<?> delegateBuilder,
256+
XdsServingStatusListener listener,
257+
FilterChainSelectorManager filterChainSelectorManager,
258+
XdsClientPoolFactory xdsClientPoolFactory,
259+
@Nullable Map<String, ?> bootstrapOverride,
260+
@Nullable Function<String, String> ldsResourceNameResolver,
261+
FilterRegistry filterRegistry,
262+
ScheduledExecutorService timeService) {
263+
this(
264+
listenerAddress,
265+
delegateBuilder,
266+
listener,
267+
filterChainSelectorManager,
268+
xdsClientPoolFactory,
269+
bootstrapOverride,
270+
ldsResourceNameResolver,
271+
filterRegistry,
272+
timeService,
273+
builder -> {});
274+
}
275+
198276
@VisibleForTesting
199277
XdsServerWrapper(
200278
String listenerAddress,
@@ -206,6 +284,31 @@ public void uncaughtException(Thread t, Throwable e) {
206284
FilterRegistry filterRegistry,
207285
ScheduledExecutorService timeService,
208286
ChannelConfigurator channelConfigurator) {
287+
this(
288+
listenerAddress,
289+
delegateBuilder,
290+
listener,
291+
filterChainSelectorManager,
292+
xdsClientPoolFactory,
293+
bootstrapOverride,
294+
null,
295+
filterRegistry,
296+
timeService,
297+
channelConfigurator);
298+
}
299+
300+
@VisibleForTesting
301+
XdsServerWrapper(
302+
String listenerAddress,
303+
ServerBuilder<?> delegateBuilder,
304+
XdsServingStatusListener listener,
305+
FilterChainSelectorManager filterChainSelectorManager,
306+
XdsClientPoolFactory xdsClientPoolFactory,
307+
@Nullable Map<String, ?> bootstrapOverride,
308+
@Nullable Function<String, String> ldsResourceNameResolver,
309+
FilterRegistry filterRegistry,
310+
ScheduledExecutorService timeService,
311+
ChannelConfigurator channelConfigurator) {
209312
this.listenerAddress = checkNotNull(listenerAddress, "listenerAddress");
210313
this.delegateBuilder = checkNotNull(delegateBuilder, "delegateBuilder");
211314
this.delegateBuilder.intercept(new ConfigApplyingInterceptor());
@@ -214,6 +317,7 @@ public void uncaughtException(Thread t, Throwable e) {
214317
= checkNotNull(filterChainSelectorManager, "filterChainSelectorManager");
215318
this.xdsClientPoolFactory = checkNotNull(xdsClientPoolFactory, "xdsClientPoolFactory");
216319
this.bootstrapOverride = bootstrapOverride;
320+
this.ldsResourceNameResolver = ldsResourceNameResolver;
217321
this.timeService = checkNotNull(timeService, "timeService");
218322
this.filterRegistry = checkNotNull(filterRegistry,"filterRegistry");
219323
this.delegate = delegateBuilder.build();
@@ -261,21 +365,28 @@ private void internalStart() {
261365
return;
262366
}
263367
xdsClient = xdsClientPool.getObject();
264-
String listenerTemplate = xdsClient.getBootstrapInfo().serverListenerResourceNameTemplate();
265-
if (listenerTemplate == null) {
266-
StatusException statusException =
267-
Status.UNAVAILABLE.withDescription(
268-
"Can only support xDS v3 with listener resource name template").asException();
269-
listener.onNotServing(statusException);
270-
initialStartFuture.set(statusException);
271-
xdsClient = xdsClientPool.returnObject(xdsClient);
272-
return;
273-
}
274-
String replacement = listenerAddress;
275-
if (listenerTemplate.startsWith(XDSTP_SCHEME)) {
276-
replacement = XdsClient.percentEncodePath(replacement);
368+
String resourceName;
369+
if (ldsResourceNameResolver != null) {
370+
resourceName = ldsResourceNameResolver.apply(listenerAddress);
371+
} else {
372+
String listenerTemplate = xdsClient.getBootstrapInfo().serverListenerResourceNameTemplate();
373+
if (listenerTemplate == null) {
374+
StatusException statusException =
375+
Status.UNAVAILABLE
376+
.withDescription("Can only support xDS v3 with listener resource name template")
377+
.asException();
378+
listener.onNotServing(statusException);
379+
initialStartFuture.set(statusException);
380+
xdsClient = xdsClientPool.returnObject(xdsClient);
381+
return;
382+
}
383+
String replacement = listenerAddress;
384+
if (listenerTemplate.startsWith(XDSTP_SCHEME)) {
385+
replacement = XdsClient.percentEncodePath(replacement);
386+
}
387+
resourceName = listenerTemplate.replaceAll("%s", replacement);
277388
}
278-
discoveryState = new DiscoveryState(listenerTemplate.replaceAll("%s", replacement));
389+
discoveryState = new DiscoveryState(resourceName);
279390
}
280391

281392
@Override

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

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -181,6 +181,50 @@ public void run() {
181181
any(SynchronizationContext.class));
182182
}
183183

184+
@Test
185+
@SuppressWarnings("unchecked")
186+
public void testBootstrap_ldsResourceNameResolver() throws Exception {
187+
Bootstrapper.BootstrapInfo b =
188+
Bootstrapper.BootstrapInfo.builder()
189+
.servers(
190+
Arrays.asList(
191+
Bootstrapper.ServerInfo.create("uri", InsecureChannelCredentials.create())))
192+
.node(EnvoyProtoData.Node.newBuilder().setId("id").build())
193+
.serverListenerResourceNameTemplate("grpc/server?udpa.resource.listening_address=%s")
194+
.build();
195+
XdsClient xdsClient = mock(XdsClient.class);
196+
XdsListenerResource listenerResource = XdsListenerResource.getInstance();
197+
when(xdsClient.getBootstrapInfo()).thenReturn(b);
198+
xdsServerWrapper =
199+
new XdsServerWrapper(
200+
"[::FFFF:129.144.52.38]:80",
201+
mockBuilder,
202+
listener,
203+
selectorManager,
204+
new FakeXdsClientPoolFactory(xdsClient),
205+
XdsServerTestHelper.RAW_BOOTSTRAP,
206+
addr -> "xdstp://resolved_name/" + addr,
207+
filterRegistry);
208+
Executors.newSingleThreadExecutor()
209+
.execute(
210+
new Runnable() {
211+
@Override
212+
public void run() {
213+
try {
214+
xdsServerWrapper.start();
215+
} catch (IOException ex) {
216+
// ignore
217+
}
218+
}
219+
});
220+
verify(xdsClient, timeout(5000))
221+
.watchXdsResource(
222+
eq(listenerResource),
223+
eq("xdstp://resolved_name/[::FFFF:129.144.52.38]:80"),
224+
any(ResourceWatcher.class),
225+
any(SynchronizationContext.class));
226+
}
227+
184228
@Test
185229
public void testBootstrap_noTemplate() throws Exception {
186230
Bootstrapper.BootstrapInfo b =

0 commit comments

Comments
 (0)