Skip to content
Merged
Show file tree
Hide file tree
Changes from 23 commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
225dc96
subscription in deps mgr
eshitachandwani Dec 26, 2025
1dd9405
remove get and add ref and other changes
eshitachandwani Dec 30, 2025
9de1462
remove create cluster ref
eshitachandwani Dec 30, 2025
49c1b2b
improve comments
eshitachandwani Jan 2, 2026
77cd824
resolve comments
eshitachandwani Jan 5, 2026
0adc4ce
resolve comments
eshitachandwani Jan 5, 2026
30d9ec7
changes
eshitachandwani Jan 7, 2026
94904ab
change ClusterSubscription api
eshitachandwani Jan 8, 2026
1aba73e
change according to design
eshitachandwani Jan 9, 2026
b656507
remove test
eshitachandwani Jan 9, 2026
ddac402
improve
eshitachandwani Jan 11, 2026
87fa9a6
Merge remote-tracking branch 'upstream/master' into change1
eshitachandwani Jan 11, 2026
730e241
Merge remote-tracking branch 'upstream/master' into change1
eshitachandwani Jan 11, 2026
d58b996
improve
eshitachandwani Jan 11, 2026
2efe5da
improve
eshitachandwani Jan 12, 2026
3ef5156
Merge branch 'master' into change1
eshitachandwani Jan 14, 2026
4cc6b7a
improve comments
eshitachandwani Jan 28, 2026
02c2812
improve test in xds_resolver_test.go
eshitachandwani Jan 28, 2026
f8edad4
improve comments
eshitachandwani Jan 28, 2026
21668a5
improve comments
eshitachandwani Jan 28, 2026
dd435dd
improve variable name
eshitachandwani Jan 29, 2026
afb58ce
remove flag
eshitachandwani Jan 29, 2026
4da7df9
update logic
eshitachandwani Jan 30, 2026
c2a7b9d
update logic
eshitachandwani Jan 31, 2026
e1008dd
change approach
eshitachandwani Feb 3, 2026
ab9440e
correct config cluster
eshitachandwani Feb 3, 2026
464ecbb
review comments
eshitachandwani Feb 4, 2026
8ef8677
review comments
eshitachandwani Feb 5, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions internal/xds/resolver/cluster_specifier_plugin_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ func (s) TestResolverClusterSpecifierPlugin(t *testing.T) {
ClusterSpecifierPluginName: "cspA",
ClusterSpecifierPluginConfig: testutils.MarshalAny(t, &wrapperspb.StringValue{Value: "anything"}),
})}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, nil, nil)

stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)

Expand Down Expand Up @@ -171,7 +171,7 @@ func (s) TestResolverClusterSpecifierPlugin(t *testing.T) {
ClusterSpecifierPluginName: "cspA",
ClusterSpecifierPluginConfig: testutils.MarshalAny(t, &wrapperspb.StringValue{Value: "changed"}),
})}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, nil, nil)

// Wait for an update from the resolver, and verify the service config.
wantSC = `
Expand Down Expand Up @@ -216,7 +216,7 @@ func (s) TestXDSResolverDelayedOnCommittedCSP(t *testing.T) {
ClusterSpecifierPluginName: "cspA",
ClusterSpecifierPluginConfig: testutils.MarshalAny(t, &wrapperspb.StringValue{Value: "anythingA"}),
})}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, nil, nil)

stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)

Expand Down Expand Up @@ -265,7 +265,7 @@ func (s) TestXDSResolverDelayedOnCommittedCSP(t *testing.T) {
ClusterSpecifierPluginName: "cspB",
ClusterSpecifierPluginConfig: testutils.MarshalAny(t, &wrapperspb.StringValue{Value: "anythingB"}),
})}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, nil, nil)

// Wait for an update from the resolver, and verify the service config.
wantSC = `
Expand Down
23 changes: 4 additions & 19 deletions internal/xds/resolver/helpers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -254,36 +254,21 @@ func setupManagementServerForTest(t *testing.T, nodeID string) (*e2e.ManagementS
return mgmtServer, listenerResourceNamesCh, routeConfigResourceNamesCh, bootstrapContents
}

// Spins up an xDS management server and configures it with a default listener
// and route configuration resource. It also sets up an xDS bootstrap
// configuration file that points to the above management server.
func configureResourcesOnManagementServer(ctx context.Context, t *testing.T, mgmtServer *e2e.ManagementServer, nodeID string, listeners []*v3listenerpb.Listener, routes []*v3routepb.RouteConfiguration) {
// Updates all resources on the given management server.
func configureResources(ctx context.Context, t *testing.T, mgmtServer *e2e.ManagementServer, nodeID string, listeners []*v3listenerpb.Listener, routes []*v3routepb.RouteConfiguration, clusters []*v3clusterpb.Cluster, endpoints []*v3endpointpb.ClusterLoadAssignment) {
resources := e2e.UpdateOptions{
NodeID: nodeID,
Listeners: listeners,
Routes: routes,
Clusters: clusters,
Endpoints: endpoints,
SkipValidation: true,
}
if err := mgmtServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}
}

// Updates all the listener, route, cluster and endpoint configuration resources
// on the given management server.
func configureAllResourcesOnManagementServer(ctx context.Context, t *testing.T, mgmtServer *e2e.ManagementServer, nodeID string, listeners []*v3listenerpb.Listener, routes []*v3routepb.RouteConfiguration, clusters []*v3clusterpb.Cluster, endpoints []*v3endpointpb.ClusterLoadAssignment) {
resources := e2e.UpdateOptions{
NodeID: nodeID,
Listeners: listeners,
Routes: routes,
Clusters: clusters,
Endpoints: endpoints,
}
if err := mgmtServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}
}

// waitForResourceNames waits for the wantNames to be pushed on to namesCh.
// Fails the test by calling t.Fatal if the context expires before that.
func waitForResourceNames(ctx context.Context, t *testing.T, namesCh chan []string, wantNames []string) {
Expand Down
6 changes: 3 additions & 3 deletions internal/xds/resolver/watch_service_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,7 @@ func (s) TestServiceWatch_ListenerPointsToNewRouteConfiguration(t *testing.T) {

// Update the management server with the new route configuration resource.
resources.Routes = append(resources.Routes, e2e.DefaultRouteConfig(newTestRouteConfigName, defaultTestServiceName, resources.Clusters[0].Name))
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, resources.Listeners, resources.Routes)
configureResources(ctx, t, mgmtServer, nodeID, resources.Listeners, resources.Routes, nil, nil)

// Ensure update from the resolver.
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(resources.Clusters[0].Name))
Expand Down Expand Up @@ -169,15 +169,15 @@ func (s) TestServiceWatch_ListenerPointsToInlineRouteConfiguration(t *testing.T)
}},
}},
}}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, resources.Listeners, nil)
configureResources(ctx, t, mgmtServer, nodeID, resources.Listeners, nil, nil, nil)

// Verify that the old route configuration is not requested anymore.
waitForResourceNames(ctx, t, routeCfgCh, []string{})
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(resources.Clusters[0].Name))

// Update listener back to contain a route configuration name.
resources.Listeners = []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, resources.Routes[0].Name)}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, resources.Listeners, resources.Routes)
configureResources(ctx, t, mgmtServer, nodeID, resources.Listeners, resources.Routes, nil, nil)

// Verify that that route configuration resource is requested.
waitForResourceNames(ctx, t, routeCfgCh, []string{resources.Routes[0].Name})
Expand Down
16 changes: 16 additions & 0 deletions internal/xds/resolver/xds_resolver.go
Comment thread
easwars marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,21 @@ func (b *xdsResolverBuilder) Build(target resolver.Target, cc resolver.ClientCon
}
r.logger = prefixLogger(r)
r.logger.Infof("Creating resolver for target: %+v", target)

dmSet := make(chan struct{})
// Schedule a callback that blocks until r.dm is set i.e xdsdepmgr.New()
// returns. This acts as a gatekeeper: even if dependency manager sends the
// updates before the xdsdepmgr.New() has a chance to return, they will be
// queued behind this blocker and processed only after initialization is
// complete.
r.serializer.TrySchedule(func(ctx context.Context) {
select {
case <-dmSet:
case <-ctx.Done():
}
})
r.dm = xdsdepmgr.New(r.ldsResourceName, opts.Authority, r.xdsClient, r)
close(dmSet)
return r, nil
}

Expand Down Expand Up @@ -329,6 +343,8 @@ func (r *xdsResolver) sendNewServiceConfig(cs stoppableConfigSelector) bool {
state := iresolver.SetConfigSelector(resolver.State{
ServiceConfig: r.cc.ParseServiceConfig(string(sc)),
}, cs)
state = xdsresource.SetXDSConfig(state, r.xdsConfig)
state = xdsdepmgr.SetXDSClusterSubscriber(state, r.dm)
if err := r.cc.UpdateState(xdsclient.SetClient(state, r.xdsClient)); err != nil {
if r.logger.V(2) {
r.logger.Infof("Channel rejected new state: %+v with error: %v", state, err)
Expand Down
122 changes: 111 additions & 11 deletions internal/xds/resolver/xds_resolver_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -232,7 +232,7 @@ func (s) TestResolverWatchCallbackAfterClose(t *testing.T) {
routes := []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(defaultTestRouteConfigName, defaultTestServiceName, defaultTestClusterName)}
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
defer cancel()
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, nil, nil)

// Wait for a discovery request for a route configuration resource.
stateCh, _, r := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, contents)
Expand Down Expand Up @@ -300,7 +300,7 @@ func (s) TestNoMatchingVirtualHost(t *testing.T) {
listener := e2e.DefaultClientListener(defaultTestServiceName, defaultTestRouteConfigName)
route := e2e.DefaultRouteConfig(defaultTestRouteConfigName, defaultTestServiceName, defaultTestClusterName)
route.VirtualHosts = nil
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{listener}, []*v3routepb.RouteConfiguration{route})
configureResources(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{listener}, []*v3routepb.RouteConfiguration{route}, nil, nil)

// Build the resolver inline (duplicating buildResolverForTarget internals)
// to avoid issues with blocked channel writes when NACKs occur.
Expand Down Expand Up @@ -369,7 +369,7 @@ func (s) TestResolverBadServiceUpdate_NACKedWithoutCache(t *testing.T) {
}},
}},
}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{lis}, nil)
configureResources(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{lis}, nil, nil, nil)

// Build the resolver inline (duplicating buildResolverForTarget internals)
// to avoid issues with blocked channel writes when NACKs occur.
Expand Down Expand Up @@ -462,7 +462,7 @@ func (s) TestResolverBadServiceUpdate_NACKedWithCache(t *testing.T) {

// Since the resource is cached, it should be received as an ambient error
// and so the RPCs should continue passing.
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{lis}, nil)
configureResources(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{lis}, nil, nil, nil)

// "Make an RPC" by invoking the config selector which should succeed by
// continuing to use the previously cached resource.
Expand Down Expand Up @@ -552,7 +552,7 @@ func (s) TestResolverGoodServiceUpdate(t *testing.T) {
// route configuration resource, as specified by the test case.
listeners := []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, defaultTestRouteConfigName)}
routes := []*v3routepb.RouteConfiguration{tt.routeConfig}
configureAllResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes, tt.clusterConfig, tt.endpointConfig)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, tt.clusterConfig, tt.endpointConfig)

stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)

Expand Down Expand Up @@ -621,7 +621,7 @@ func (s) TestResolverRequestHash(t *testing.T) {
}}
cluster := []*v3clusterpb.Cluster{e2e.DefaultCluster(defaultTestClusterName, defaultTestEndpointName, e2e.SecurityLevelNone)}
endpoints := []*v3endpointpb.ClusterLoadAssignment{e2e.DefaultEndpoint(defaultTestEndpointName, defaultTestHostname, defaultTestPort)}
configureAllResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes, cluster, endpoints)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, cluster, endpoints)

// Build the resolver and read the config selector out of it.
stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)
Expand Down Expand Up @@ -908,7 +908,7 @@ func (s) TestResolverMaxStreamDuration(t *testing.T) {
e2e.DefaultEndpoint("endpoint_B", defaultTestHostname, defaultTestPort),
e2e.DefaultEndpoint("endpoint_C", defaultTestHostname, defaultTestPort),
}
configureAllResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes, cluster, endpoints)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, cluster, endpoints)

// Read the update pushed by the resolver to the ClientConn.
cs := verifyUpdateFromResolver(ctx, t, stateCh, "")
Expand Down Expand Up @@ -1084,7 +1084,7 @@ func (s) TestResolverMultipleLDSUpdates(t *testing.T) {
// Configure the management server with a listener resource, but no route
// configuration resource.
listeners := []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, defaultTestRouteConfigName)}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, nil)
configureResources(ctx, t, mgmtServer, nodeID, listeners, nil, nil, nil)

stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)

Expand Down Expand Up @@ -1118,7 +1118,7 @@ func (s) TestResolverMultipleLDSUpdates(t *testing.T) {
}},
}},
}}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, nil)
configureResources(ctx, t, mgmtServer, nodeID, listeners, nil, nil, nil)

// Ensure that there is no update from the resolver.
verifyNoUpdateFromResolver(ctx, t, stateCh)
Expand Down Expand Up @@ -1157,7 +1157,7 @@ func (s) TestResolverWRR(t *testing.T) {
e2e.DefaultEndpoint("endpoint_A", defaultTestHostname, defaultTestPort),
e2e.DefaultEndpoint("endpoint_B", defaultTestHostname, defaultTestPort),
}
configureAllResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes, clusters, endpoints)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, clusters, endpoints)

// Read the update pushed by the resolver to the ClientConn.
cs := verifyUpdateFromResolver(ctx, t, stateCh, "")
Expand Down Expand Up @@ -1251,7 +1251,7 @@ func (s) TestConfigSelector_FailureCases(t *testing.T) {

// Update the management server with a listener resource that
// contains inline route configuration.
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{test.listener}, nil)
configureResources(ctx, t, mgmtServer, nodeID, []*v3listenerpb.Listener{test.listener}, nil, nil, nil)

// Ensure that the resolver pushes a state update to the channel.
cs := verifyUpdateFromResolver(ctx, t, stateCh, "")
Expand Down Expand Up @@ -1418,3 +1418,103 @@ func (s) TestResolver_AutoHostRewrite(t *testing.T) {
})
}
}

// TestResolverKeepWatchOpen_ActiveRPCs tests that the dependency manager keeps
// a cluster watch open when there are active RPCs using that cluster, even if
// the cluster is no longer referenced by the current route configuration.
func (s) TestResolverKeepWatchOpen_ActiveRPCs(t *testing.T) {
t.Skip("Will be enabled when all the watchers have shifted to dependency manager")

clusterA := "cluster-A"
clusterB := "cluster-B"

ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
defer cancel()
cdsResourceRequestedCh := make(chan []string, 2)
mgmtServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{
OnStreamRequest: func(_ int64, req *v3discoverypb.DiscoveryRequest) error {
if req.GetTypeUrl() != version.V3ClusterURL {
return nil
}
if len(req.GetResourceNames()) > 0 {
select {
case cdsResourceRequestedCh <- req.GetResourceNames():
case <-ctx.Done():
}
}
return nil
},
AllowResourceSubset: true,
})

nodeID := uuid.New().String()
bc := e2e.DefaultBootstrapContents(t, nodeID, mgmtServer.Address)

// Configure initial resources: Route -> ClusterA.
listeners := []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, defaultTestRouteConfigName)}
Comment thread
easwars marked this conversation as resolved.
routes := []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(defaultTestRouteConfigName, defaultTestServiceName, clusterA)}
clusters := []*v3clusterpb.Cluster{
e2e.DefaultCluster(clusterA, "endpoint-A", e2e.SecurityLevelNone),
e2e.DefaultCluster(clusterB, "endpoint-B", e2e.SecurityLevelNone),
}
endpoints := []*v3endpointpb.ClusterLoadAssignment{
e2e.DefaultEndpoint("endpoint-A", "localhost", []uint32{8080}),
e2e.DefaultEndpoint("endpoint-B", "localhost", []uint32{8081}),
}

configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, clusters, endpoints)

stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)

cs := verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(clusterA))

// Start RPC (Ref Counts ClusterA).
res, err := cs.SelectConfig(iresolver.RPCInfo{Context: ctx, Method: "/service/method"})
if err != nil {
t.Fatalf("cs.SelectConfig(): %v", err)
}

// Switch Configuration to ClusterB.
routes[0] = e2e.DefaultRouteConfig(defaultTestRouteConfigName, defaultTestServiceName, clusterB)
configureResources(ctx, t, mgmtServer, nodeID, listeners, routes, clusters, endpoints)

// Resolver should request BOTH A (due to active RPC) and B (due to new
// config).
wantNames := []string{clusterA, clusterB}
waitForResourceNames(ctx, t, cdsResourceRequestedCh, wantNames)

// Verify Service Config has both clusters.
const wantServiceRaw = `{
"loadBalancingConfig": [{
"xds_cluster_manager_experimental": {
"children": {
"cluster:cluster-A": {
"childPolicy": [{"cds_experimental": {"cluster": "cluster-A"}}]
},
"cluster:cluster-B": {
"childPolicy": [{"cds_experimental": {"cluster": "cluster-B"}}]
}
}
}
}]
}`
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceRaw)

// Finish RPC (Drops Ref to ClusterA).
res.OnCommitted()

// ONLY cluster B should be requested now that there are no references to
// cluster A.
wantNames = []string{clusterB}
waitForResourceNames(ctx, t, cdsResourceRequestedCh, wantNames)

// ServiceConfig update should also contain only cluster B.
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(clusterB))
Comment thread
easwars marked this conversation as resolved.

// Verify that RPCs pass again.
res, err = cs.SelectConfig(iresolver.RPCInfo{Context: ctx, Method: "/service/method"})
if err != nil {
t.Fatalf("cs.SelectConfig(): %v", err)
}
res.OnCommitted()
}
Loading