From ed8e278c4023b8c6995e0de7a0e46b86b986f661 Mon Sep 17 00:00:00 2001 From: Ran Date: Thu, 23 Jul 2026 16:36:03 +0000 Subject: [PATCH 1/2] Supports injecting custom ResourceNameResolvers --- .../java/io/grpc/xds/XdsServerBuilder.java | 43 +++++- .../java/io/grpc/xds/XdsServerWrapper.java | 143 ++++++++++++++++-- .../io/grpc/xds/XdsServerWrapperTest.java | 44 ++++++ 3 files changed, 212 insertions(+), 18 deletions(-) diff --git a/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java b/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java index 1c0eb3cd024..b1845a4abe4 100644 --- a/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java +++ b/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java @@ -37,10 +37,14 @@ import io.grpc.netty.InternalProtocolNegotiator; import io.grpc.netty.NettyServerBuilder; import io.grpc.xds.FilterChainMatchingProtocolNegotiators.FilterChainMatchingNegotiatorServerFactory; +import java.net.InetSocketAddress; +import java.net.SocketAddress; import java.util.Map; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Function; import java.util.logging.Logger; +import javax.annotation.Nullable; /** * A version of {@link ServerBuilder} to create xDS managed servers. @@ -57,6 +61,7 @@ public final class XdsServerBuilder extends ForwardingServerBuilder bootstrapOverride; + @Nullable private Function ldsResourceNameResolver; private long drainGraceTime = 10; private TimeUnit drainGraceTimeUnit = TimeUnit.MINUTES; private ChannelConfigurator channelConfigurator = builder -> { }; @@ -134,6 +139,23 @@ public static XdsServerBuilder forPort(int port, ServerCredentials serverCredent return new XdsServerBuilder(nettyDelegate, port); } + /** Creates a gRPC server builder for the given address. */ + public static XdsServerBuilder forAddress( + SocketAddress address, ServerCredentials serverCredentials) { + checkNotNull(serverCredentials, "serverCredentials"); + InternalProtocolNegotiator.ServerFactory originalNegotiatorFactory = + InternalNettyServerCredentials.toNegotiator(serverCredentials); + ServerCredentials wrappedCredentials = + InternalNettyServerCredentials.create( + new FilterChainMatchingNegotiatorServerFactory(originalNegotiatorFactory)); + NettyServerBuilder nettyDelegate = NettyServerBuilder.forAddress(address, wrappedCredentials); + int port = 0; + if (address instanceof InetSocketAddress inetSocketAddress) { + port = inetSocketAddress.getPort(); + } + return new XdsServerBuilder(nettyDelegate, port); + } + @Override public Server build() { checkState(isServerBuilt.compareAndSet(false, true), "Server already built!"); @@ -144,11 +166,28 @@ public Server build() { builder.set(ATTR_DRAIN_GRACE_NANOS, drainGraceTimeUnit.toNanos(drainGraceTime)); } InternalNettyServerBuilder.eagAttributes(delegate, builder.build()); - return new XdsServerWrapper("0.0.0.0:" + port, delegate, xdsServingStatusListener, - filterChainSelectorManager, xdsClientPoolFactory, bootstrapOverride, filterRegistry, + return new XdsServerWrapper( + "0.0.0.0:" + port, + delegate, + xdsServingStatusListener, + filterChainSelectorManager, + xdsClientPoolFactory, + bootstrapOverride, + ldsResourceNameResolver, + filterRegistry, this.channelConfigurator); } + /** + * Provides a function that takes the listening address and returns the LDS resource name. When + * provided, this overrides the server_listener_resource_name_template in the bootstrap. + */ + public XdsServerBuilder ldsResourceNameResolver( + Function ldsResourceNameResolver) { + this.ldsResourceNameResolver = checkNotNull(ldsResourceNameResolver, "ldsResourceNameResolver"); + return this; + } + @VisibleForTesting XdsServerBuilder xdsClientPoolFactory(XdsClientPoolFactory xdsClientPoolFactory) { this.xdsClientPoolFactory = checkNotNull(xdsClientPoolFactory, "xdsClientPoolFactory"); diff --git a/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java b/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java index 8aeae7f4988..56f30d37e5d 100644 --- a/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java +++ b/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java @@ -78,6 +78,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Function; import java.util.logging.Level; import java.util.logging.Logger; import javax.annotation.Nullable; @@ -108,6 +109,7 @@ public void uncaughtException(Thread t, Throwable e) { private final ThreadSafeRandom random = ThreadSafeRandomImpl.instance; private final XdsClientPoolFactory xdsClientPoolFactory; private final @Nullable Map bootstrapOverride; + private final @Nullable Function ldsResourceNameResolver; private final XdsServingStatusListener listener; private final FilterChainSelectorManager filterChainSelectorManager; private final AtomicBoolean started = new AtomicBoolean(false); @@ -139,6 +141,7 @@ public void uncaughtException(Thread t, Throwable e) { FilterChainSelectorManager filterChainSelectorManager, XdsClientPoolFactory xdsClientPoolFactory, @Nullable Map bootstrapOverride, + @Nullable Function ldsResourceNameResolver, FilterRegistry filterRegistry, ChannelConfigurator channelConfigurator) { this( @@ -148,6 +151,30 @@ public void uncaughtException(Thread t, Throwable e) { filterChainSelectorManager, xdsClientPoolFactory, bootstrapOverride, + ldsResourceNameResolver, + filterRegistry, + SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE), + channelConfigurator); + sharedTimeService = true; + } + + XdsServerWrapper( + String listenerAddress, + ServerBuilder delegateBuilder, + XdsServingStatusListener listener, + FilterChainSelectorManager filterChainSelectorManager, + XdsClientPoolFactory xdsClientPoolFactory, + @Nullable Map bootstrapOverride, + FilterRegistry filterRegistry, + ChannelConfigurator channelConfigurator) { + this( + listenerAddress, + delegateBuilder, + listener, + filterChainSelectorManager, + xdsClientPoolFactory, + bootstrapOverride, + null, filterRegistry, SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE), channelConfigurator); @@ -169,8 +196,34 @@ public void uncaughtException(Thread t, Throwable e) { filterChainSelectorManager, xdsClientPoolFactory, bootstrapOverride, + null, filterRegistry, - builder -> { }); + SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE), + builder -> {}); + sharedTimeService = true; + } + + XdsServerWrapper( + String listenerAddress, + ServerBuilder delegateBuilder, + XdsServingStatusListener listener, + FilterChainSelectorManager filterChainSelectorManager, + XdsClientPoolFactory xdsClientPoolFactory, + @Nullable Map bootstrapOverride, + @Nullable Function ldsResourceNameResolver, + FilterRegistry filterRegistry) { + this( + listenerAddress, + delegateBuilder, + listener, + filterChainSelectorManager, + xdsClientPoolFactory, + bootstrapOverride, + ldsResourceNameResolver, + filterRegistry, + SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE), + builder -> {}); + sharedTimeService = true; } @VisibleForTesting @@ -190,9 +243,34 @@ public void uncaughtException(Thread t, Throwable e) { filterChainSelectorManager, xdsClientPoolFactory, bootstrapOverride, + null, + filterRegistry, + timeService, + builder -> {}); + } + + @VisibleForTesting + XdsServerWrapper( + String listenerAddress, + ServerBuilder delegateBuilder, + XdsServingStatusListener listener, + FilterChainSelectorManager filterChainSelectorManager, + XdsClientPoolFactory xdsClientPoolFactory, + @Nullable Map bootstrapOverride, + @Nullable Function ldsResourceNameResolver, + FilterRegistry filterRegistry, + ScheduledExecutorService timeService) { + this( + listenerAddress, + delegateBuilder, + listener, + filterChainSelectorManager, + xdsClientPoolFactory, + bootstrapOverride, + ldsResourceNameResolver, filterRegistry, timeService, - builder -> { }); + builder -> {}); } @VisibleForTesting @@ -206,6 +284,31 @@ public void uncaughtException(Thread t, Throwable e) { FilterRegistry filterRegistry, ScheduledExecutorService timeService, ChannelConfigurator channelConfigurator) { + this( + listenerAddress, + delegateBuilder, + listener, + filterChainSelectorManager, + xdsClientPoolFactory, + bootstrapOverride, + null, + filterRegistry, + timeService, + channelConfigurator); + } + + @VisibleForTesting + XdsServerWrapper( + String listenerAddress, + ServerBuilder delegateBuilder, + XdsServingStatusListener listener, + FilterChainSelectorManager filterChainSelectorManager, + XdsClientPoolFactory xdsClientPoolFactory, + @Nullable Map bootstrapOverride, + @Nullable Function ldsResourceNameResolver, + FilterRegistry filterRegistry, + ScheduledExecutorService timeService, + ChannelConfigurator channelConfigurator) { this.listenerAddress = checkNotNull(listenerAddress, "listenerAddress"); this.delegateBuilder = checkNotNull(delegateBuilder, "delegateBuilder"); this.delegateBuilder.intercept(new ConfigApplyingInterceptor()); @@ -214,6 +317,7 @@ public void uncaughtException(Thread t, Throwable e) { = checkNotNull(filterChainSelectorManager, "filterChainSelectorManager"); this.xdsClientPoolFactory = checkNotNull(xdsClientPoolFactory, "xdsClientPoolFactory"); this.bootstrapOverride = bootstrapOverride; + this.ldsResourceNameResolver = ldsResourceNameResolver; this.timeService = checkNotNull(timeService, "timeService"); this.filterRegistry = checkNotNull(filterRegistry,"filterRegistry"); this.delegate = delegateBuilder.build(); @@ -261,21 +365,28 @@ private void internalStart() { return; } xdsClient = xdsClientPool.getObject(); - String listenerTemplate = xdsClient.getBootstrapInfo().serverListenerResourceNameTemplate(); - if (listenerTemplate == null) { - StatusException statusException = - Status.UNAVAILABLE.withDescription( - "Can only support xDS v3 with listener resource name template").asException(); - listener.onNotServing(statusException); - initialStartFuture.set(statusException); - xdsClient = xdsClientPool.returnObject(xdsClient); - return; - } - String replacement = listenerAddress; - if (listenerTemplate.startsWith(XDSTP_SCHEME)) { - replacement = XdsClient.percentEncodePath(replacement); + String resourceName; + if (ldsResourceNameResolver != null) { + resourceName = ldsResourceNameResolver.apply(listenerAddress); + } else { + String listenerTemplate = xdsClient.getBootstrapInfo().serverListenerResourceNameTemplate(); + if (listenerTemplate == null) { + StatusException statusException = + Status.UNAVAILABLE + .withDescription("Can only support xDS v3 with listener resource name template") + .asException(); + listener.onNotServing(statusException); + initialStartFuture.set(statusException); + xdsClient = xdsClientPool.returnObject(xdsClient); + return; + } + String replacement = listenerAddress; + if (listenerTemplate.startsWith(XDSTP_SCHEME)) { + replacement = XdsClient.percentEncodePath(replacement); + } + resourceName = listenerTemplate.replaceAll("%s", replacement); } - discoveryState = new DiscoveryState(listenerTemplate.replaceAll("%s", replacement)); + discoveryState = new DiscoveryState(resourceName); } @Override diff --git a/xds/src/test/java/io/grpc/xds/XdsServerWrapperTest.java b/xds/src/test/java/io/grpc/xds/XdsServerWrapperTest.java index cc26195b039..47ac32cdc8a 100644 --- a/xds/src/test/java/io/grpc/xds/XdsServerWrapperTest.java +++ b/xds/src/test/java/io/grpc/xds/XdsServerWrapperTest.java @@ -181,6 +181,50 @@ public void run() { any(SynchronizationContext.class)); } + @Test + @SuppressWarnings("unchecked") + public void testBootstrap_ldsResourceNameResolver() throws Exception { + Bootstrapper.BootstrapInfo b = + Bootstrapper.BootstrapInfo.builder() + .servers( + Arrays.asList( + Bootstrapper.ServerInfo.create("uri", InsecureChannelCredentials.create()))) + .node(EnvoyProtoData.Node.newBuilder().setId("id").build()) + .serverListenerResourceNameTemplate("grpc/server?udpa.resource.listening_address=%s") + .build(); + XdsClient xdsClient = mock(XdsClient.class); + XdsListenerResource listenerResource = XdsListenerResource.getInstance(); + when(xdsClient.getBootstrapInfo()).thenReturn(b); + xdsServerWrapper = + new XdsServerWrapper( + "[::FFFF:129.144.52.38]:80", + mockBuilder, + listener, + selectorManager, + new FakeXdsClientPoolFactory(xdsClient), + XdsServerTestHelper.RAW_BOOTSTRAP, + addr -> "xdstp://resolved_name/" + addr, + filterRegistry); + Executors.newSingleThreadExecutor() + .execute( + new Runnable() { + @Override + public void run() { + try { + xdsServerWrapper.start(); + } catch (IOException ex) { + // ignore + } + } + }); + verify(xdsClient, timeout(5000)) + .watchXdsResource( + eq(listenerResource), + eq("xdstp://resolved_name/[::FFFF:129.144.52.38]:80"), + any(ResourceWatcher.class), + any(SynchronizationContext.class)); + } + @Test public void testBootstrap_noTemplate() throws Exception { Bootstrapper.BootstrapInfo b = From f9c6f35c23243552b36c89ed4bf4485b4ea269f0 Mon Sep 17 00:00:00 2001 From: Ran Date: Thu, 23 Jul 2026 17:15:40 +0000 Subject: [PATCH 2/2] stop using pattern matching --- xds/src/main/java/io/grpc/xds/XdsServerBuilder.java | 3 ++- xds/src/main/java/io/grpc/xds/XdsServerWrapper.java | 8 ++++---- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java b/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java index b1845a4abe4..7bef1725af2 100644 --- a/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java +++ b/xds/src/main/java/io/grpc/xds/XdsServerBuilder.java @@ -150,7 +150,8 @@ public static XdsServerBuilder forAddress( new FilterChainMatchingNegotiatorServerFactory(originalNegotiatorFactory)); NettyServerBuilder nettyDelegate = NettyServerBuilder.forAddress(address, wrappedCredentials); int port = 0; - if (address instanceof InetSocketAddress inetSocketAddress) { + if (address instanceof InetSocketAddress) { + InetSocketAddress inetSocketAddress = (InetSocketAddress) address; port = inetSocketAddress.getPort(); } return new XdsServerBuilder(nettyDelegate, port); diff --git a/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java b/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java index 56f30d37e5d..dffdf2c7476 100644 --- a/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java +++ b/xds/src/main/java/io/grpc/xds/XdsServerWrapper.java @@ -199,7 +199,7 @@ public void uncaughtException(Thread t, Throwable e) { null, filterRegistry, SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE), - builder -> {}); + builder -> { }); sharedTimeService = true; } @@ -222,7 +222,7 @@ public void uncaughtException(Thread t, Throwable e) { ldsResourceNameResolver, filterRegistry, SharedResourceHolder.get(GrpcUtil.TIMER_SERVICE), - builder -> {}); + builder -> { }); sharedTimeService = true; } @@ -246,7 +246,7 @@ public void uncaughtException(Thread t, Throwable e) { null, filterRegistry, timeService, - builder -> {}); + builder -> { }); } @VisibleForTesting @@ -270,7 +270,7 @@ public void uncaughtException(Thread t, Throwable e) { ldsResourceNameResolver, filterRegistry, timeService, - builder -> {}); + builder -> { }); } @VisibleForTesting