From 533dafbade540399eea69ab7ebece7d29f2e6413 Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Wed, 15 Jul 2026 11:13:59 +0000 Subject: [PATCH 01/10] A75: Implement the new LB policy topology for the non-aggregate clusters only --- .../cdsbalancer/aggregate_cluster_test.go | 88 ++---- .../xds/balancer/cdsbalancer/cdsbalancer.go | 52 +++- .../balancer/cdsbalancer/cdsbalancer_test.go | 282 +++++++++++------- .../xds/balancer/cdsbalancer/configbuilder.go | 108 ++++++- .../cdsbalancer/configbuilder_test.go | 16 +- .../cdsbalancer/e2e_test/balancer_test.go | 68 ++--- 6 files changed, 392 insertions(+), 222 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go b/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go index 9cb040d5c3ba..b56bab30186d 100644 --- a/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go +++ b/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go @@ -90,8 +90,8 @@ func verifyDNSResolution(ctx context.Context, t *testing.T, dnsTargetCh chan res // Tests the case where the cluster resource requested is a leaf cluster. The // management server sends two updates for the same leaf cluster resource. The -// test verifies that the load balancing configuration pushed to the priority LB -// policy contains the expected discovery mechanism corresponding to the leaf +// test verifies that the load balancing configuration pushed to the top-level LB +// policy contains the expected configuration corresponding to the leaf // cluster, on both occasions. func (s) TestAggregateClusterSuccess_LeafNode(t *testing.T) { tests := []struct { @@ -105,43 +105,21 @@ func (s) TestAggregateClusterSuccess_LeafNode(t *testing.T) { name: "eds", firstClusterResource: e2e.DefaultCluster(clusterName, serviceName, e2e.SecurityLevelNone), secondClusterResource: e2e.DefaultCluster(clusterName, serviceName+"-new", e2e.SecurityLevelNone), - wantFirstChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: createPriorityConfig(clusterName), - IgnoreReresolutionRequests: true, - }, - }, - Priorities: []string{"priority-0-0"}, - }, - wantSecondChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-1-0": { - Config: createPriorityConfig(clusterName), - IgnoreReresolutionRequests: true, - }, - }, - Priorities: []string{"priority-1-0"}, - }, + wantFirstChildCfg: createSingleClusterConfig(clusterName, "priority-0-0", true), + wantSecondChildCfg: createSingleClusterConfig(clusterName, "priority-1-0", true), }, { name: "dns", firstClusterResource: makeLogicalDNSClusterResource(clusterName, "dns_host", uint32(port)), secondClusterResource: makeLogicalDNSClusterResource(clusterName, "dns_host_new", uint32(port)), - wantFirstChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{"priority-0": {Config: createPriorityConfig(clusterName)}}, - Priorities: []string{"priority-0"}, - }, - wantSecondChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{"priority-1": {Config: createPriorityConfig(clusterName)}}, - Priorities: []string{"priority-1"}, - }, + wantFirstChildCfg: createSingleClusterConfig(clusterName, "priority-0", false), + wantSecondChildCfg: createSingleClusterConfig(clusterName, "priority-1", false), }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { - lbCfgCh, _, _, _ := registerWrappedPriorityPolicy(t) + lbCfgCh := registerWrappedOutlierDetectionPolicy(t) mgmtServer, nodeID, _ := setupWithManagementServer(t, nil, nil) // Push the first cluster resource through the management server and @@ -281,11 +259,11 @@ func (s) TestAggregateClusterSuccess_ThenUpdateChildClusters(t *testing.T) { // the priority LB policy contains the discovery mechanisms for both child // clusters. The test then updates the root cluster resource requested by the // cds LB policy to a leaf cluster of type EDS and verifies the load balancing -// configuration pushed to the priority LB policy contains a single discovery -// mechanism. +// configuration pushed to the outlier detection LB policy contains the leaf cluster config. func (s) TestAggregateClusterSuccess_ThenChangeRootToEDS(t *testing.T) { dnsTargetCh, dnsR := setupDNS(t) lbCfgCh, _, _, _ := registerWrappedPriorityPolicy(t) + odCfgCh := registerWrappedOutlierDetectionPolicy(t) mgmtServer, nodeID, _ := setupWithManagementServer(t, nil, nil) // Configure the management server with the aggregate cluster resource @@ -332,21 +310,17 @@ func (s) TestAggregateClusterSuccess_ThenChangeRootToEDS(t *testing.T) { }, Endpoints: []*v3endpointpb.ClusterLoadAssignment{e2e.DefaultEndpoint(serviceName, host, []uint32{port})}, } + // Drain odCfgCh before sending non-aggregate update, as aggregate children pushed outlier configs into odCfgCh. + for len(odCfgCh) > 0 { + <-odCfgCh + } if err := mgmtServer.Update(ctx, resources); err != nil { t.Fatal(err) } // Since the service name of the EDS cluster remains same, same priority name // is used. - wantChildCfg = &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: createPriorityConfig(clusterName), - IgnoreReresolutionRequests: true, - }, - }, - Priorities: []string{"priority-0-0"}, - } - if err := compareLoadBalancingConfig(ctx, lbCfgCh, wantChildCfg); err != nil { + wantSingleChildCfg := createSingleClusterConfig(clusterName, "priority-0-0", true) + if err := compareLoadBalancingConfig(ctx, odCfgCh, wantSingleChildCfg); err != nil { t.Fatal(err) } } @@ -354,10 +328,11 @@ func (s) TestAggregateClusterSuccess_ThenChangeRootToEDS(t *testing.T) { // Tests the case where a requested cluster resource switches between being a // leaf and an aggregate cluster pointing to an EDS and LogicalDNS child // cluster. In each of these cases, the test verifies that the load balancing -// configuration pushed to the priority LB policy contains the expected config. +// configuration pushed to the top-level child policy contains the expected config. func (s) TestAggregatedClusterSuccess_SwitchBetweenLeafAndAggregate(t *testing.T) { dnsTargetCh, dnsR := setupDNS(t) lbCfgCh, _, _, _ := registerWrappedPriorityPolicy(t) + odCfgCh := registerWrappedOutlierDetectionPolicy(t) mgmtServer, nodeID, _ := setupWithManagementServer(t, nil, nil) // Start off with the requested cluster being a leaf EDS cluster. @@ -373,16 +348,8 @@ func (s) TestAggregatedClusterSuccess_SwitchBetweenLeafAndAggregate(t *testing.T if err := mgmtServer.Update(ctx, resources); err != nil { t.Fatal(err) } - wantChildCfg := &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: createPriorityConfig(clusterName), - IgnoreReresolutionRequests: true, - }, - }, - Priorities: []string{"priority-0-0"}, - } - if err := compareLoadBalancingConfig(ctx, lbCfgCh, wantChildCfg); err != nil { + wantSingleChildCfg := createSingleClusterConfig(clusterName, "priority-0-0", true) + if err := compareLoadBalancingConfig(ctx, odCfgCh, wantSingleChildCfg); err != nil { t.Fatal(err) } @@ -404,7 +371,7 @@ func (s) TestAggregatedClusterSuccess_SwitchBetweenLeafAndAggregate(t *testing.T } verifyDNSResolution(ctx, t, dnsTargetCh, dnsR, dnsHostName, dnsPort) - wantChildCfg = &priority.LBConfig{ + wantChildCfg := &priority.LBConfig{ Children: map[string]*priority.Child{ "priority-0-0": { Config: createPriorityConfig(edsClusterName), @@ -426,19 +393,14 @@ func (s) TestAggregatedClusterSuccess_SwitchBetweenLeafAndAggregate(t *testing.T Clusters: []*v3clusterpb.Cluster{e2e.DefaultCluster(clusterName, serviceName, e2e.SecurityLevelNone)}, Endpoints: []*v3endpointpb.ClusterLoadAssignment{e2e.DefaultEndpoint(serviceName, host, []uint32{port})}, } + // Drain odCfgCh before sending non-aggregate update, as aggregate children pushed outlier configs into odCfgCh. + for len(odCfgCh) > 0 { + <-odCfgCh + } if err := mgmtServer.Update(ctx, resources); err != nil { t.Fatal(err) } - wantChildCfg = &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: createPriorityConfig(clusterName), - IgnoreReresolutionRequests: true, - }, - }, - Priorities: []string{"priority-0-0"}, - } - if err := compareLoadBalancingConfig(ctx, lbCfgCh, wantChildCfg); err != nil { + if err := compareLoadBalancingConfig(ctx, odCfgCh, wantSingleChildCfg); err != nil { t.Fatal(err) } } diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 08b6b451510d..6bcdd308bcda 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -41,16 +41,20 @@ import ( const cdsName = "cds_experimental" var ( - // newChildBalancer is a helper function to build a new priority balancer - // and will be overridden in unittests. - newChildBalancer = func(cc balancer.ClientConn, opts balancer.BuildOptions) (balancer.Balancer, error) { - builder := balancer.Get(priority.Name) + // newChildBalancer is a helper function to build a new child balancer and its + // config parser, and will be overridden in unittests. + newChildBalancer = func(name string, cc balancer.ClientConn, opts balancer.BuildOptions) (balancer.Balancer, balancer.ConfigParser, error) { + builder := balancer.Get(name) if builder == nil { - return nil, fmt.Errorf("xds: no balancer builder with name %v", priority.Name) + return nil, nil, fmt.Errorf("xds: no balancer builder with name %v", name) } - // We directly pass the parent clientConn to the underlying priority + parser, ok := builder.(balancer.ConfigParser) + if !ok { + return nil, nil, fmt.Errorf("xds: balancer builder for %v does not implement ConfigParser", name) + } + // We directly pass the parent clientConn to the underlying child // balancer because the cdsBalancer does not deal with subConns. - return builder.Build(cc, opts), nil + return builder.Build(cc, opts), parser, nil } ) @@ -134,6 +138,7 @@ type cdsBalancer struct { // protect access to these fields. xdsClient xdsclient.XDSClient childLB balancer.Balancer // Child policy, built upon resolution of the cluster graph. + childLBName string // Name of the child policy. clusterConfigs map[string]*xdsresource.ClusterResult // Cluster name to the last received result for that cluster. priorityConfigs map[string]*priorityConfig // Hostname to priority config for that leaf cluster. lbCfg *lbConfig // Current load balancing configuration. @@ -267,18 +272,45 @@ func (b *cdsBalancer) handleClusterUpdate() error { // A child policy is created if one doesn't already exist. The newly built // configuration is then pushed to the child policy. func (b *cdsBalancer) updateChildConfig() error { + clusterName := b.lbCfg.ClusterName + clusterConfig := b.clusterConfigs[clusterName].Config + isAggregate := clusterConfig.Cluster.ClusterType == xdsresource.ClusterTypeAggregate + + var topLBName string + if isAggregate { + topLBName = priority.Name + } else { + topLBName = outlierdetection.Name + } + + if b.childLB != nil && b.childLBName != topLBName { + b.childLB.Close() + b.childLB = nil + } + if b.childLB == nil { - childLB, err := newChildBalancer(b.cc, b.bOpts) + childLB, parser, err := newChildBalancer(topLBName, b.cc, b.bOpts) if err != nil { - return fmt.Errorf("failed to create child policy of type %s: %v", priority.Name, err) + return fmt.Errorf("failed to create child policy of type %s: %v", topLBName, err) } b.childLB = childLB + b.childLBName = topLBName + b.childConfigParser = parser } - childCfgBytes, endpoints, err := buildPriorityConfigJSON(b.priorities, &b.xdsLBPolicy) + var childCfgBytes []byte + var endpoints []resolver.Endpoint + var err error + + if isAggregate { + childCfgBytes, endpoints, err = buildAggregateClusterPriorityConfigJSON(b.priorities, &b.xdsLBPolicy) + } else { + childCfgBytes, endpoints, err = buildSingleClusterConfigJSON(b.priorities[0], &b.xdsLBPolicy) + } if err != nil { return fmt.Errorf("failed to build child policy config: %v", err) } + childCfg, err := b.childConfigParser.ParseConfig(childCfgBytes) if err != nil { return fmt.Errorf("failed to parse child policy config. This should never happen because the config was generated: %v", err) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go b/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go index 7e1a9b2201fa..77a787c68aa9 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go @@ -152,11 +152,19 @@ func registerWrappedPriorityPolicy(t *testing.T) (chan serviceconfig.LoadBalanci }, ExitIdle: func(bd *stub.BalancerData) { bd.ChildBalancer.ExitIdle() - close(exitIdleCh) + select { + case <-exitIdleCh: + default: + close(exitIdleCh) + } }, Close: func(bd *stub.BalancerData) { bd.ChildBalancer.Close() - close(closeCh) + select { + case <-closeCh: + default: + close(closeCh) + } }, }) t.Cleanup(func() { balancer.Register(priorityBuilder) }) @@ -164,6 +172,35 @@ func registerWrappedPriorityPolicy(t *testing.T) (chan serviceconfig.LoadBalanci return lbCfgCh, resolverErrCh, exitIdleCh, closeCh } +func registerWrappedOutlierDetectionPolicy(t *testing.T) chan serviceconfig.LoadBalancingConfig { + odBuilder := balancer.Get(outlierdetection.Name) + internal.BalancerUnregister(odBuilder.Name()) + + lbCfgCh := make(chan serviceconfig.LoadBalancingConfig, 1) + + stub.Register(outlierdetection.Name, stub.BalancerFuncs{ + Init: func(bd *stub.BalancerData) { + bd.ChildBalancer = odBuilder.Build(bd.ClientConn, bd.BuildOptions) + }, + ParseConfig: func(lbCfg json.RawMessage) (serviceconfig.LoadBalancingConfig, error) { + return odBuilder.(balancer.ConfigParser).ParseConfig(lbCfg) + }, + UpdateClientConnState: func(bd *stub.BalancerData, ccs balancer.ClientConnState) error { + select { + case lbCfgCh <- ccs.BalancerConfig: + default: + } + return bd.ChildBalancer.UpdateClientConnState(ccs) + }, + Close: func(bd *stub.BalancerData) { + bd.ChildBalancer.Close() + }, + }) + t.Cleanup(func() { balancer.Register(odBuilder) }) + + return lbCfgCh +} + // setupDNS unregisters the DNS resolver and registers a manual resolver for the // same scheme. This allows the test to fake the DNS resolution by supplying the // addresses of the test backends. @@ -236,19 +273,29 @@ func compareLoadBalancingConfig(ctx context.Context, lbCfgCh chan serviceconfig. if err != nil { return fmt.Errorf("failed to marshal expected child config to JSON: %v", err) } - select { - case lbCfg := <-lbCfgCh: - gotJSON, err := json.Marshal(lbCfg) - if err != nil { - return fmt.Errorf("failed to marshal received LB config into JSON: %v", err) - } - if diff := cmp.Diff(wantJSON, gotJSON); diff != "" { - return fmt.Errorf("child policy received unexpected diff in config (-want +got):\n%s", diff) + var lastErr error + // Loop over updates received on lbCfgCh to consume intermediate configuration + // updates (e.g. pushed as individual discovery mechanisms resolve) until a + // configuration matching wantChildCfg is received, or the context deadline expires. + for { + select { + case lbCfg := <-lbCfgCh: + gotJSON, err := json.Marshal(lbCfg) + if err != nil { + return fmt.Errorf("failed to marshal received LB config into JSON: %v", err) + } + if diff := cmp.Diff(wantJSON, gotJSON); diff != "" { + lastErr = fmt.Errorf("child policy received unexpected diff in config (-want +got):\n%s", diff) + continue + } + return nil + case <-ctx.Done(): + if lastErr != nil { + return lastErr + } + return fmt.Errorf("timeout when waiting for child policy to receive its configuration") } - case <-ctx.Done(): - return fmt.Errorf("timeout when waiting for child policy to receive its configuration") } - return nil } func verifyRPCError(gotErr error, wantCode codes.Code, wantErr, wantNodeID string) error { @@ -290,6 +337,41 @@ func createPriorityConfig(cluster string) *iserviceconfig.BalancerConfig { } } +// createSingleClusterConfig returns the expected LoadBalancingConfig tree for a +// single (non-aggregate) cluster under gRFC A75 topology: +// outlier_detection -> cluster_impl -> priority -> wrr_locality -> round_robin. +func createSingleClusterConfig(cluster string, pName string, ignoreReresolution bool) *outlierdetection.LBConfig { + return &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: cluster, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + pName: { + Config: &iserviceconfig.BalancerConfig{ + Name: wrrlocality.Name, + Config: &wrrlocality.LBConfig{ + ChildPolicy: &iserviceconfig.BalancerConfig{Name: roundrobin.Name}, + }, + }, + IgnoreReresolutionRequests: ignoreReresolution, + }, + }, + Priorities: []string{pName}, + }, + }, + }, + }, + } +} + // Tests the case where a configuration with an empty cluster name is pushed to // the CDS LB policy. Verifies that ErrBadResolverState is returned. func (s) TestConfigurationUpdate_EmptyCluster(t *testing.T) { @@ -414,32 +496,32 @@ func (s) TestClusterUpdate_Success(t *testing.T) { } return c }(), - wantChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: &iserviceconfig.BalancerConfig{ - Name: outlierdetection.Name, - Config: &outlierdetection.LBConfig{ - Interval: iserviceconfig.Duration(10 * time.Second), // default interval - BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), - MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), - MaxEjectionPercent: 10, - ChildPolicy: &iserviceconfig.BalancerConfig{ - Name: clusterimpl.Name, - Config: &clusterimpl.LBConfig{ - Cluster: clusterName, - ChildPolicy: &iserviceconfig.BalancerConfig{ + wantChildCfg: &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), // default interval + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: clusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-0-0": { + Config: &iserviceconfig.BalancerConfig{ Name: wrrlocality.Name, Config: &wrrlocality.LBConfig{ChildPolicy: &iserviceconfig.BalancerConfig{Name: roundrobin.Name}}, }, + IgnoreReresolutionRequests: true, }, }, + Priorities: []string{"priority-0-0"}, }, }, - IgnoreReresolutionRequests: true, }, }, - Priorities: []string{"priority-0-0"}, }, }, { @@ -459,35 +541,35 @@ func (s) TestClusterUpdate_Success(t *testing.T) { } return c }(), - wantChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: &iserviceconfig.BalancerConfig{ - Name: outlierdetection.Name, - Config: &outlierdetection.LBConfig{ - Interval: iserviceconfig.Duration(10 * time.Second), // default interval - BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), - MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), - MaxEjectionPercent: 10, - ChildPolicy: &iserviceconfig.BalancerConfig{ - Name: clusterimpl.Name, - Config: &clusterimpl.LBConfig{ - Cluster: clusterName, - ChildPolicy: &iserviceconfig.BalancerConfig{ + wantChildCfg: &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), // default interval + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: clusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-0-0": { + Config: &iserviceconfig.BalancerConfig{ Name: ringhash.Name, Config: &iringhash.LBConfig{ MinRingSize: 100, MaxRingSize: 1000, }, }, + IgnoreReresolutionRequests: true, }, }, + Priorities: []string{"priority-0-0"}, }, }, - IgnoreReresolutionRequests: true, }, }, - Priorities: []string{"priority-0-0"}, }, }, { @@ -502,41 +584,41 @@ func (s) TestClusterUpdate_Success(t *testing.T) { c.OutlierDetection = &v3clusterpb.OutlierDetection{} return c }(), - wantChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: &iserviceconfig.BalancerConfig{ - Name: outlierdetection.Name, - Config: &outlierdetection.LBConfig{ - Interval: iserviceconfig.Duration(10 * time.Second), // default interval - BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), - MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), - MaxEjectionPercent: 10, - SuccessRateEjection: &outlierdetection.SuccessRateEjection{ - StdevFactor: 1900, - EnforcementPercentage: 100, - MinimumHosts: 5, - RequestVolume: 100, - }, - ChildPolicy: &iserviceconfig.BalancerConfig{ - Name: clusterimpl.Name, - Config: &clusterimpl.LBConfig{ - Cluster: clusterName, - ChildPolicy: &iserviceconfig.BalancerConfig{ + wantChildCfg: &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), // default interval + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + SuccessRateEjection: &outlierdetection.SuccessRateEjection{ + StdevFactor: 1900, + EnforcementPercentage: 100, + MinimumHosts: 5, + RequestVolume: 100, + }, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: clusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-0-0": { + Config: &iserviceconfig.BalancerConfig{ Name: ringhash.Name, Config: &iringhash.LBConfig{ MinRingSize: 1024, // default sizes MaxRingSize: 4096, }, }, + IgnoreReresolutionRequests: true, }, }, + Priorities: []string{"priority-0-0"}, }, }, - IgnoreReresolutionRequests: true, }, }, - Priorities: []string{"priority-0-0"}, }, }, { @@ -564,54 +646,54 @@ func (s) TestClusterUpdate_Success(t *testing.T) { } return c }(), - wantChildCfg: &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: &iserviceconfig.BalancerConfig{ - Name: outlierdetection.Name, - Config: &outlierdetection.LBConfig{ - Interval: iserviceconfig.Duration(10 * time.Second), // default interval - BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), - MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), - MaxEjectionPercent: 10, - SuccessRateEjection: &outlierdetection.SuccessRateEjection{ - StdevFactor: 1900, - EnforcementPercentage: 100, - MinimumHosts: 5, - RequestVolume: 100, - }, - FailurePercentageEjection: &outlierdetection.FailurePercentageEjection{ - Threshold: 85, - EnforcementPercentage: 5, - MinimumHosts: 5, - RequestVolume: 50, - }, - ChildPolicy: &iserviceconfig.BalancerConfig{ - Name: clusterimpl.Name, - Config: &clusterimpl.LBConfig{ - Cluster: clusterName, - ChildPolicy: &iserviceconfig.BalancerConfig{ + wantChildCfg: &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), // default interval + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + SuccessRateEjection: &outlierdetection.SuccessRateEjection{ + StdevFactor: 1900, + EnforcementPercentage: 100, + MinimumHosts: 5, + RequestVolume: 100, + }, + FailurePercentageEjection: &outlierdetection.FailurePercentageEjection{ + Threshold: 85, + EnforcementPercentage: 5, + MinimumHosts: 5, + RequestVolume: 50, + }, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: clusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-0-0": { + Config: &iserviceconfig.BalancerConfig{ Name: ringhash.Name, Config: &iringhash.LBConfig{ MinRingSize: 1024, // default sizes MaxRingSize: 4096, }, }, + IgnoreReresolutionRequests: true, }, }, + Priorities: []string{"priority-0-0"}, }, }, - IgnoreReresolutionRequests: true, }, }, - Priorities: []string{"priority-0-0"}, }, }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { - lbCfgCh, _, _, _ := registerWrappedPriorityPolicy(t) + lbCfgCh := registerWrappedOutlierDetectionPolicy(t) mgmtServer, nodeID, _ := setupWithManagementServer(t, nil, nil) ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) diff --git a/internal/xds/balancer/cdsbalancer/configbuilder.go b/internal/xds/balancer/cdsbalancer/configbuilder.go index 6efeeb14b65e..4babc133645d 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder.go @@ -77,10 +77,106 @@ func hostName(clusterName string, update xdsresource.ClusterUpdate) string { } } -// buildPriorityConfigJSON builds balancer config for the passed in -// priorities. +// buildSingleClusterConfigJSON builds the balancer config for a single +// (non-aggregate) cluster according to gRFC A75. // -// The built tree of balancers (see test for the output struct). +// The built tree of balancers: +// +// ┌─────────────────┐ +// │outlier_detection│ +// └────────┬────────┘ +// │ +// ┌────────▼───────┐ +// │ cluster_impl │ +// └────────┬───────┘ +// │ +// ┌───▼────┐ +// │priority│ +// └┬──────┬┘ +// │ │ +// ┌──────────▼─┐ ┌─▼──────────┐ +// │xDSLBPolicy │ │xDSLBPolicy │ (Locality and Endpoint picking layer) +// └────────────┘ └────────────┘ +func buildSingleClusterConfigJSON(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { + odCfg, endpoints, err := buildSingleClusterConfig(p, xdsLBPolicy) + if err != nil { + return nil, nil, err + } + ret, err := json.Marshal(odCfg) + if err != nil { + return nil, nil, fmt.Errorf("failed to marshal built single cluster config: %v", err) + } + return ret, endpoints, nil +} + +func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*outlierdetection.LBConfig, []resolver.Endpoint, error) { + clusterUpdate := p.clusterConfig.Cluster + priorityLBConfig := &priority.LBConfig{ + Children: make(map[string]*priority.Child), + } + var retEndpoints []resolver.Endpoint + + switch clusterUpdate.ClusterType { + case xdsresource.ClusterTypeEDS: + priorities := [][]xdsresource.Locality{{}} + edsUpdate := p.clusterConfig.EndpointConfig.EDSUpdate + if len(edsUpdate.Localities) != 0 { + priorities = groupLocalitiesByPriority(edsUpdate.Localities) + } + priorityNames := p.childNameGen.generate(priorities) + priorityLBConfig.Priorities = priorityNames + + for i, pName := range priorityNames { + priorityLocalities := priorities[i] + _, endpoints, err := priorityLocalitiesToClusterImpl(priorityLocalities, pName, *p.clusterConfig.Cluster, xdsLBPolicy) + if err != nil { + return nil, nil, err + } + retEndpoints = append(retEndpoints, endpoints...) + priorityLBConfig.Children[pName] = &priority.Child{ + Config: xdsLBPolicy, + IgnoreReresolutionRequests: true, + } + } + case xdsresource.ClusterTypeLogicalDNS: + pName := fmt.Sprintf("priority-%v", p.childNameGen.prefix) + priorityLBConfig.Priorities = []string{pName} + endpoints := p.clusterConfig.EndpointConfig.DNSEndpoints.Endpoints + var retEndpoint resolver.Endpoint + for _, e := range endpoints { + retEndpoint.Addresses = append(retEndpoint.Addresses, e.Addresses...) + } + localityStr := xdsinternal.LocalityString(clients.Locality{}) + retEndpoint = hostname.Set(hierarchy.SetInEndpoint(retEndpoint, []string{pName, localityStr}), clusterUpdate.DNSHostName) + retEndpoint = wrrlocality.SetAddrInfo(retEndpoint, wrrlocality.AddrInfo{LocalityWeight: 1}) + retEndpoints = append(retEndpoints, retEndpoint) + priorityLBConfig.Children[pName] = &priority.Child{ + Config: xdsLBPolicy, + IgnoreReresolutionRequests: false, + } + } + + ciCfg := &clusterimpl.LBConfig{ + Cluster: clusterUpdate.ClusterName, + ChildPolicy: &internalserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: priorityLBConfig, + }, + } + + odCfg := p.outlierDetection + odCfg.ChildPolicy = &internalserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: ciCfg, + } + + return &odCfg, retEndpoints, nil +} + +// buildAggregateClusterPriorityConfigJSON builds balancer config for the passed in +// priorities (legacy / aggregate cluster tree). +// +// The built tree of balancers: // // ┌────────┐ // │priority│ @@ -93,8 +189,8 @@ func hostName(clusterName string, update xdsresource.ClusterUpdate) string { // ┌──────▼─────┐ ┌─────▼──────┐ // │xDSLBPolicy │ │xDSLBPolicy │ (Locality and Endpoint picking layer) // └────────────┘ └────────────┘ -func buildPriorityConfigJSON(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { - pc, endpoints, err := buildPriorityConfig(priorities, xdsLBPolicy) +func buildAggregateClusterPriorityConfigJSON(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { + pc, endpoints, err := buildAggregateClusterPriorityConfig(priorities, xdsLBPolicy) if err != nil { return nil, nil, fmt.Errorf("failed to build priority config: %v", err) } @@ -105,7 +201,7 @@ func buildPriorityConfigJSON(priorities []*priorityConfig, xdsLBPolicy *internal return ret, endpoints, nil } -func buildPriorityConfig(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*priority.LBConfig, []resolver.Endpoint, error) { +func buildAggregateClusterPriorityConfig(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*priority.LBConfig, []resolver.Endpoint, error) { var ( retConfig = &priority.LBConfig{Children: make(map[string]*priority.Child)} retEndpoints []resolver.Endpoint diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 282e2a49f6ea..784074d2e155 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -130,9 +130,9 @@ func makeLocality(localityIdx int, localityWeight, priority uint32, endpointCoun } } -// TestBuildPriorityConfigJSON is a sanity check that the built balancer config -// can be parsed. The behavior test is covered by TestBuildPriorityConfig. -func (s) TestBuildPriorityConfigJSON(t *testing.T) { +// TestBuildAggregateClusterPriorityConfigJSON is a sanity check that the built balancer config +// can be parsed. The behavior test is covered by TestBuildAggregateClusterPriorityConfig. +func (s) TestBuildAggregateClusterPriorityConfigJSON(t *testing.T) { testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ URI: "trafficdirector.googleapis.com:443", ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, @@ -141,7 +141,7 @@ func (s) TestBuildPriorityConfigJSON(t *testing.T) { t.Fatalf("Failed to create LRS server config for testing: %v", err) } - gotConfig, _, err := buildPriorityConfigJSON([]*priorityConfig{ + gotConfig, _, err := buildAggregateClusterPriorityConfigJSON([]*priorityConfig{ { clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ @@ -182,7 +182,7 @@ func (s) TestBuildPriorityConfigJSON(t *testing.T) { }, }, nil) if err != nil { - t.Fatalf("buildPriorityConfigJSON(...) failed: %v", err) + t.Fatalf("buildAggregateClusterPriorityConfigJSON(...) failed: %v", err) } var prettyGot bytes.Buffer @@ -198,11 +198,11 @@ func (s) TestBuildPriorityConfigJSON(t *testing.T) { } } -// TestBuildPriorityConfig tests the priority config generation. Each top level +// TestBuildAggregateClusterPriorityConfig tests the priority config generation. Each top level // balancer per priority should be an Outlier Detection balancer, with a Cluster // Impl Balancer as a child. -func (s) TestBuildPriorityConfig(t *testing.T) { - gotConfig, _, _ := buildPriorityConfig([]*priorityConfig{ +func (s) TestBuildAggregateClusterPriorityConfig(t *testing.T) { + gotConfig, _, _ := buildAggregateClusterPriorityConfig([]*priorityConfig{ { // EDS - OD config should be the top level for both of the EDS // priorities balancer This EDS priority will have multiple sub diff --git a/internal/xds/balancer/cdsbalancer/e2e_test/balancer_test.go b/internal/xds/balancer/cdsbalancer/e2e_test/balancer_test.go index 6b7df2f26d07..64bab9171a47 100644 --- a/internal/xds/balancer/cdsbalancer/e2e_test/balancer_test.go +++ b/internal/xds/balancer/cdsbalancer/e2e_test/balancer_test.go @@ -300,21 +300,21 @@ func (s) TestErrorFromParentLB_ResourceNotFound(t *testing.T) { } // Test verifies that when the received Cluster resource contains outlier -// detection configuration, the LB config pushed to the priority policy contains -// the appropriate configuration for the outlier detection LB policy. +// detection configuration, the LB config pushed to the outlier detection policy +// contains the appropriate configuration for the outlier detection LB policy. func (s) TestOutlierDetectionConfigPropagationToChildPolicy(t *testing.T) { - // Unregister the priority balancer builder for the duration of this test, + // Unregister the outlier detection balancer builder for the duration of this test, // and register a policy under the same name that makes the LB config // pushed to it available to the test. - priorityBuilder := balancer.Get(priority.Name) - internal.BalancerUnregister(priorityBuilder.Name()) + odBuilder := balancer.Get(outlierdetection.Name) + internal.BalancerUnregister(odBuilder.Name()) lbCfgCh := make(chan serviceconfig.LoadBalancingConfig, 1) - stub.Register(priority.Name, stub.BalancerFuncs{ + stub.Register(outlierdetection.Name, stub.BalancerFuncs{ Init: func(bd *stub.BalancerData) { - bd.ChildBalancer = priorityBuilder.Build(bd.ClientConn, bd.BuildOptions) + bd.ChildBalancer = odBuilder.Build(bd.ClientConn, bd.BuildOptions) }, ParseConfig: func(lbCfg json.RawMessage) (serviceconfig.LoadBalancingConfig, error) { - return priorityBuilder.(balancer.ConfigParser).ParseConfig(lbCfg) + return odBuilder.(balancer.ConfigParser).ParseConfig(lbCfg) }, UpdateClientConnState: func(bd *stub.BalancerData, ccs balancer.ClientConnState) error { select { @@ -327,7 +327,7 @@ func (s) TestOutlierDetectionConfigPropagationToChildPolicy(t *testing.T) { bd.ChildBalancer.Close() }, }) - defer balancer.Register(priorityBuilder) + defer balancer.Register(odBuilder) managementServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{}) @@ -368,29 +368,27 @@ func (s) TestOutlierDetectionConfigPropagationToChildPolicy(t *testing.T) { _, cleanup := setupAndDial(t, bootstrapContents) defer cleanup() - // The priority configuration generated should have Outlier Detection as a - // direct child due to Outlier Detection being turned on. - wantCfg := &priority.LBConfig{ - Children: map[string]*priority.Child{ - "priority-0-0": { - Config: &iserviceconfig.BalancerConfig{ - Name: outlierdetection.Name, - Config: &outlierdetection.LBConfig{ - Interval: iserviceconfig.Duration(10 * time.Second), // default interval - BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), - MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), - MaxEjectionPercent: 10, - SuccessRateEjection: &outlierdetection.SuccessRateEjection{ - StdevFactor: 2000, - EnforcementPercentage: 50, - MinimumHosts: 10, - RequestVolume: 50, - }, - ChildPolicy: &iserviceconfig.BalancerConfig{ - Name: clusterimpl.Name, - Config: &clusterimpl.LBConfig{ - Cluster: clusterName, - ChildPolicy: &iserviceconfig.BalancerConfig{ + wantCfg := &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), // default interval + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + SuccessRateEjection: &outlierdetection.SuccessRateEjection{ + StdevFactor: 2000, + EnforcementPercentage: 50, + MinimumHosts: 10, + RequestVolume: 50, + }, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: clusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-0-0": { + Config: &iserviceconfig.BalancerConfig{ Name: wrrlocality.Name, Config: &wrrlocality.LBConfig{ ChildPolicy: &iserviceconfig.BalancerConfig{ @@ -398,19 +396,19 @@ func (s) TestOutlierDetectionConfigPropagationToChildPolicy(t *testing.T) { }, }, }, + IgnoreReresolutionRequests: true, }, }, + Priorities: []string{"priority-0-0"}, }, }, - IgnoreReresolutionRequests: true, }, }, - Priorities: []string{"priority-0-0"}, } select { case lbCfg := <-lbCfgCh: - gotCfg := lbCfg.(*priority.LBConfig) + gotCfg := lbCfg.(*outlierdetection.LBConfig) if diff := cmp.Diff(wantCfg, gotCfg); diff != "" { t.Fatalf("Child policy received unexpected diff in config (-want +got):\n%s", diff) } From 9b3c703668bee4b435fc9ed127b573382c4814b0 Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Fri, 17 Jul 2026 11:58:38 +0000 Subject: [PATCH 02/10] A75: Add missing test cases --- .../cdsbalancer/configbuilder_test.go | 364 ++++++++++++++++++ 1 file changed, 364 insertions(+) diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 784074d2e155..7e93f940be16 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -317,6 +317,370 @@ func testEndpointForDNS(endpoints []resolver.Endpoint, localityWeight uint32, pa return retEndpoint } +func (s) TestBuildSingleClusterConfigJSON(t *testing.T) { + testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ + URI: "trafficdirector.googleapis.com:443", + ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, + }) + if err != nil { + t.Fatalf("Failed to create LRS server config for testing: %v", err) + } + + gotConfig, _, err := buildSingleClusterConfigJSON(&priorityConfig{ + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterName: testClusterName, + ClusterType: xdsresource.ClusterTypeEDS, + EDSServiceName: testEDSServiceName, + MaxRequests: newUint32(testMaxRequests), + LRSServerConfig: testLRSServerConfig, + }, + EndpointConfig: &xdsresource.EndpointConfig{ + EDSUpdate: &xdsresource.EndpointsUpdate{ + Drops: []xdsresource.OverloadDropConfig{{ + Category: testDropCategory, + Numerator: testDropOverMillion, + Denominator: million, + }}, + Localities: []xdsresource.Locality{ + makeLocality(0, 20, 0, 2), + makeLocality(1, 80, 0, 2), + makeLocality(2, 20, 1, 2), + makeLocality(3, 80, 1, 2), + }, + }, + }, + }, + outlierDetection: noopODCfg, + childNameGen: newNameGenerator(0), + }, &iserviceconfig.BalancerConfig{Name: roundrobin.Name}) + if err != nil { + t.Fatalf("buildSingleClusterConfigJSON(...) failed: %v", err) + } + + var prettyGot bytes.Buffer + if err := json.Indent(&prettyGot, gotConfig, ">>> ", " "); err != nil { + t.Fatalf("json.Indent() failed: %v", err) + } + t.Log(prettyGot.String()) + + odB := balancer.Get(outlierdetection.Name) + if _, err = odB.(balancer.ConfigParser).ParseConfig(gotConfig); err != nil { + t.Fatalf("ParseConfig(%+v) failed: %v", gotConfig, err) + } +} + +func (s) TestBuildSingleClusterConfig_DNS(t *testing.T) { + for _, tt := range []struct { + name string + endpoints []resolver.Endpoint + xdsLBPolicy *iserviceconfig.BalancerConfig + }{ + { + name: "one_endpoint_one_address", + endpoints: []resolver.Endpoint{{Addresses: []resolver.Address{{Addr: "addr-0-0"}}}}, + xdsLBPolicy: &iserviceconfig.BalancerConfig{Name: pickfirst.Name}, + }, + { + name: "one_endpoint_multiple_addresses", + endpoints: []resolver.Endpoint{{Addresses: []resolver.Address{ + {Addr: "addr-0-0"}, + {Addr: "addr-0-1"}, + }}}, + xdsLBPolicy: &iserviceconfig.BalancerConfig{Name: wrrlocality.Name}, + }, + { + name: "multiple_endpoints_one_address_each", + endpoints: []resolver.Endpoint{ + {Addresses: []resolver.Address{{Addr: "addr-0-0"}}}, + {Addresses: []resolver.Address{{Addr: "addr-0-1"}}}, + }, + xdsLBPolicy: &iserviceconfig.BalancerConfig{Name: roundrobin.Name}, + }, + { + name: "multiple_endpoints_multiple_addresses", + endpoints: []resolver.Endpoint{ + {Addresses: []resolver.Address{ + {Addr: "addr-0-0"}, + {Addr: "addr-0-1"}, + }}, + {Addresses: []resolver.Address{ + {Addr: "addr-1-0"}, + {Addr: "addr-1-1"}, + }}, + }, + xdsLBPolicy: &iserviceconfig.BalancerConfig{Name: roundrobin.Name}, + }, + } { + t.Run(tt.name, func(t *testing.T) { + gotODConfig, gotEndpoints, err := buildSingleClusterConfig( + &priorityConfig{ + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterName: testClusterName2, + ClusterType: xdsresource.ClusterTypeLogicalDNS, + }, + EndpointConfig: &xdsresource.EndpointConfig{DNSEndpoints: &xdsresource.DNSUpdate{Endpoints: tt.endpoints}}, + }, + outlierDetection: noopODCfg, + childNameGen: newNameGenerator(3), + }, + tt.xdsLBPolicy) + if err != nil { + t.Fatalf("buildSingleClusterConfig() failed: %v", err) + } + + wantODConfig := &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: testClusterName2, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-3": { + Config: tt.xdsLBPolicy, + IgnoreReresolutionRequests: false, + }, + }, + Priorities: []string{"priority-3"}, + }, + }, + }, + }, + } + if diff := cmp.Diff(wantODConfig, gotODConfig); diff != "" { + t.Errorf("buildSingleClusterConfig() config diff (-want +got) %v", diff) + } + + wantEndpoints := []resolver.Endpoint{testEndpointForDNS(tt.endpoints, 1, []string{"priority-3", xdsinternal.LocalityString(clients.Locality{})})} + if diff := cmp.Diff(wantEndpoints, gotEndpoints, endpointCmpOpts); diff != "" { + t.Errorf("buildSingleClusterConfig() endpoints diff (-want +got) %v", diff) + } + }) + } +} + +func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t *testing.T) { + testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, false) + + testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ + URI: "trafficdirector.googleapis.com:443", + ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, + }) + if err != nil { + t.Fatalf("Failed to create LRS server config for testing: %v", err) + } + + // Create test localities with 2 endpoints each. + // Localities are passed in shuffled order to verify sorting by priority. + loc0 := makeLocality(0, 20, 0, 2) + loc1 := makeLocality(1, 80, 0, 2) + loc2 := makeLocality(2, 20, 1, 2) + loc3 := makeLocality(3, 80, 1, 2) + + gotODConfig, gotEndpoints, err := buildSingleClusterConfig( + &priorityConfig{ + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterName: testClusterName, + ClusterType: xdsresource.ClusterTypeEDS, + EDSServiceName: testEDSServiceName, + LRSServerConfig: testLRSServerConfig, + MaxRequests: newUint32(testMaxRequests), + }, + EndpointConfig: &xdsresource.EndpointConfig{ + EDSUpdate: &xdsresource.EndpointsUpdate{ + Drops: []xdsresource.OverloadDropConfig{{ + Category: testDropCategory, + Numerator: testDropOverMillion, + Denominator: million, + }}, + Localities: []xdsresource.Locality{ + makeLocality(3, 80, 1, 2), + makeLocality(1, 80, 0, 2), + makeLocality(2, 20, 1, 2), + makeLocality(0, 20, 0, 2), + }, + }, + }, + }, + outlierDetection: noopODCfg, + childNameGen: newNameGenerator(2), + }, + nil, + ) + if err != nil { + t.Fatalf("buildSingleClusterConfig() failed: %v", err) + } + + wantODConfig := &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: testClusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-2-0": { + IgnoreReresolutionRequests: true, + }, + "priority-2-1": { + IgnoreReresolutionRequests: true, + }, + }, + Priorities: []string{"priority-2-0", "priority-2-1"}, + }, + }, + }, + }, + } + // Endpoint weight is the product of locality weight and endpoint weight. + wantEndpoints := []resolver.Endpoint{ + testEndpointWithAttrs(loc0.Endpoints[0].ResolverEndpoint, 20, 20*1, "priority-2-0", &loc0.ID), + testEndpointWithAttrs(loc0.Endpoints[1].ResolverEndpoint, 20, 20*1, "priority-2-0", &loc0.ID), + testEndpointWithAttrs(loc1.Endpoints[0].ResolverEndpoint, 80, 80*1, "priority-2-0", &loc1.ID), + testEndpointWithAttrs(loc1.Endpoints[1].ResolverEndpoint, 80, 80*1, "priority-2-0", &loc1.ID), + testEndpointWithAttrs(loc2.Endpoints[0].ResolverEndpoint, 20, 20*1, "priority-2-1", &loc2.ID), + testEndpointWithAttrs(loc2.Endpoints[1].ResolverEndpoint, 20, 20*1, "priority-2-1", &loc2.ID), + testEndpointWithAttrs(loc3.Endpoints[0].ResolverEndpoint, 80, 80*1, "priority-2-1", &loc3.ID), + testEndpointWithAttrs(loc3.Endpoints[1].ResolverEndpoint, 80, 80*1, "priority-2-1", &loc3.ID), + } + + if diff := cmp.Diff(wantODConfig, gotODConfig); diff != "" { + t.Errorf("buildSingleClusterConfig() config diff (-want +got) %v", diff) + } + if diff := cmp.Diff(wantEndpoints, gotEndpoints, endpointCmpOpts); diff != "" { + t.Errorf("buildSingleClusterConfig() endpoints diff (-want +got) %v", diff) + } +} + +func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t *testing.T) { + testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, true) + + testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ + URI: "trafficdirector.googleapis.com:443", + ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, + }) + if err != nil { + t.Fatalf("Failed to create LRS server config for testing: %v", err) + } + + // Create test localities with 2 endpoints each. + // Localities are passed in shuffled order to verify sorting by priority. + loc0 := makeLocality(0, 20, 0, 2) + loc1 := makeLocality(1, 80, 0, 2) + loc2 := makeLocality(2, 20, 1, 2) + loc3 := makeLocality(3, 80, 1, 2) + + gotODConfig, gotEndpoints, err := buildSingleClusterConfig( + &priorityConfig{ + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterName: testClusterName, + ClusterType: xdsresource.ClusterTypeEDS, + EDSServiceName: testEDSServiceName, + LRSServerConfig: testLRSServerConfig, + MaxRequests: newUint32(testMaxRequests), + }, + EndpointConfig: &xdsresource.EndpointConfig{ + EDSUpdate: &xdsresource.EndpointsUpdate{ + Drops: []xdsresource.OverloadDropConfig{{ + Category: testDropCategory, + Numerator: testDropOverMillion, + Denominator: million, + }}, + Localities: []xdsresource.Locality{ + makeLocality(3, 80, 1, 2), + makeLocality(1, 80, 0, 2), + makeLocality(2, 20, 1, 2), + makeLocality(0, 20, 0, 2), + }, + }, + }, + }, + outlierDetection: noopODCfg, + childNameGen: newNameGenerator(2), + }, + nil, + ) + if err != nil { + t.Fatalf("buildSingleClusterConfig() failed: %v", err) + } + + wantODConfig := &outlierdetection.LBConfig{ + Interval: iserviceconfig.Duration(10 * time.Second), + BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), + MaxEjectionTime: iserviceconfig.Duration(300 * time.Second), + MaxEjectionPercent: 10, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: clusterimpl.Name, + Config: &clusterimpl.LBConfig{ + Cluster: testClusterName, + ChildPolicy: &iserviceconfig.BalancerConfig{ + Name: priority.Name, + Config: &priority.LBConfig{ + Children: map[string]*priority.Child{ + "priority-2-0": { + IgnoreReresolutionRequests: true, + }, + "priority-2-1": { + IgnoreReresolutionRequests: true, + }, + }, + Priorities: []string{"priority-2-0", "priority-2-1"}, + }, + }, + }, + }, + } + // Endpoints weights are the product of normalized locality weight and + // endpoint weight, represented as a fixed-point number in uQ1.31 format. + // Locality weights are normalized as: + // P1: locality 3: 80 / (100) = 0.8 + // P0: locality 1: 80 / (100) = 0.8 + // P1: locality 2: 20 / (100) = 0.2 + // P0: locality 0: 20 / (100) = 0.2 + // In fixed-point uQ1.31 format, the weights are: + // locality 3: 0.8 * 2^31 = 1717986918 + // locality 1: 0.8 * 2^31 = 1717986918 + // locality 2: 0.2 * 2^31 = 429496729 + // locality 0: 0.2 * 2^31 = 429496729 + // + // There are two endpoints in each locality, each with weight 1. So, their + // normalized weights are 0.5 each. And the final endpoint weights are a + // product of their locality weights and 0.5, which turns out to be either + // 1717986918 * 0.5 = 858993459, or, + // 429496729 * 0.5 = 214748364 + wantEndpoints := []resolver.Endpoint{ + testEndpointWithAttrs(loc0.Endpoints[0].ResolverEndpoint, 20, 214748364, "priority-2-0", &loc0.ID), + testEndpointWithAttrs(loc0.Endpoints[1].ResolverEndpoint, 20, 214748364, "priority-2-0", &loc0.ID), + testEndpointWithAttrs(loc1.Endpoints[0].ResolverEndpoint, 80, 858993459, "priority-2-0", &loc1.ID), + testEndpointWithAttrs(loc1.Endpoints[1].ResolverEndpoint, 80, 858993459, "priority-2-0", &loc1.ID), + testEndpointWithAttrs(loc2.Endpoints[0].ResolverEndpoint, 20, 214748364, "priority-2-1", &loc2.ID), + testEndpointWithAttrs(loc2.Endpoints[1].ResolverEndpoint, 20, 214748364, "priority-2-1", &loc2.ID), + testEndpointWithAttrs(loc3.Endpoints[0].ResolverEndpoint, 80, 858993459, "priority-2-1", &loc3.ID), + testEndpointWithAttrs(loc3.Endpoints[1].ResolverEndpoint, 80, 858993459, "priority-2-1", &loc3.ID), + } + + if diff := cmp.Diff(wantODConfig, gotODConfig); diff != "" { + t.Errorf("buildSingleClusterConfig() config diff (-want +got) %v", diff) + } + if diff := cmp.Diff(wantEndpoints, gotEndpoints, endpointCmpOpts); diff != "" { + t.Errorf("buildSingleClusterConfig() endpoints diff (-want +got) %v", diff) + } +} + func (s) TestBuildClusterImplConfigForDNS(t *testing.T) { for _, tt := range []struct { name string From fbed93a2c516f848e7e45ec395d3a2e94236dced Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Mon, 20 Jul 2026 07:00:30 +0000 Subject: [PATCH 03/10] Address Gemini review - fix nil checks --- .../xds/balancer/cdsbalancer/cdsbalancer.go | 12 ++- .../xds/balancer/cdsbalancer/configbuilder.go | 9 ++ .../cdsbalancer/configbuilder_test.go | 82 +++++++++++++++++++ 3 files changed, 102 insertions(+), 1 deletion(-) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 6bcdd308bcda..d84df15eac20 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -273,7 +273,14 @@ func (b *cdsBalancer) handleClusterUpdate() error { // configuration is then pushed to the child policy. func (b *cdsBalancer) updateChildConfig() error { clusterName := b.lbCfg.ClusterName - clusterConfig := b.clusterConfigs[clusterName].Config + clusterResult := b.clusterConfigs[clusterName] + if clusterResult == nil { + return fmt.Errorf("cluster result for %q is missing", clusterName) + } + clusterConfig := clusterResult.Config + if clusterConfig.Cluster == nil { + return fmt.Errorf("cluster update is missing for cluster %q", clusterName) + } isAggregate := clusterConfig.Cluster.ClusterType == xdsresource.ClusterTypeAggregate var topLBName string @@ -305,6 +312,9 @@ func (b *cdsBalancer) updateChildConfig() error { if isAggregate { childCfgBytes, endpoints, err = buildAggregateClusterPriorityConfigJSON(b.priorities, &b.xdsLBPolicy) } else { + if len(b.priorities) == 0 { + return fmt.Errorf("no priorities configured for non-aggregate cluster %q", clusterName) + } childCfgBytes, endpoints, err = buildSingleClusterConfigJSON(b.priorities[0], &b.xdsLBPolicy) } if err != nil { diff --git a/internal/xds/balancer/cdsbalancer/configbuilder.go b/internal/xds/balancer/cdsbalancer/configbuilder.go index 4babc133645d..322203390c24 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder.go @@ -110,6 +110,9 @@ func buildSingleClusterConfigJSON(p *priorityConfig, xdsLBPolicy *internalservic } func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*outlierdetection.LBConfig, []resolver.Endpoint, error) { + if p == nil || p.clusterConfig == nil || p.clusterConfig.Cluster == nil || p.childNameGen == nil { + return nil, nil, fmt.Errorf("cluster config is missing") + } clusterUpdate := p.clusterConfig.Cluster priorityLBConfig := &priority.LBConfig{ Children: make(map[string]*priority.Child), @@ -119,6 +122,9 @@ func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalservicecon switch clusterUpdate.ClusterType { case xdsresource.ClusterTypeEDS: priorities := [][]xdsresource.Locality{{}} + if p.clusterConfig.EndpointConfig == nil || p.clusterConfig.EndpointConfig.EDSUpdate == nil { + return nil, nil, fmt.Errorf("EDS update is missing for cluster %q", clusterUpdate.ClusterName) + } edsUpdate := p.clusterConfig.EndpointConfig.EDSUpdate if len(edsUpdate.Localities) != 0 { priorities = groupLocalitiesByPriority(edsUpdate.Localities) @@ -141,6 +147,9 @@ func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalservicecon case xdsresource.ClusterTypeLogicalDNS: pName := fmt.Sprintf("priority-%v", p.childNameGen.prefix) priorityLBConfig.Priorities = []string{pName} + if p.clusterConfig.EndpointConfig == nil || p.clusterConfig.EndpointConfig.DNSEndpoints == nil { + return nil, nil, fmt.Errorf("DNS endpoints are missing for cluster %q", clusterUpdate.ClusterName) + } endpoints := p.clusterConfig.EndpointConfig.DNSEndpoints.Endpoints var retEndpoint resolver.Endpoint for _, e := range endpoints { diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 7e93f940be16..431e610480ba 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -681,6 +681,88 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t * } } +func (s) TestBuildSingleClusterConfig_NilChecks(t *testing.T) { + tests := []struct { + name string + p *priorityConfig + }{ + { + name: "nil priorityConfig", + p: nil, + }, + { + name: "priorityConfig with nil fields", + p: &priorityConfig{}, + }, + { + name: "nil clusterConfig", + p: &priorityConfig{childNameGen: newNameGenerator(0)}, + }, + { + name: "nil childNameGen", + p: &priorityConfig{ + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeEDS}, + }, + }, + }, + { + name: "nil Cluster in clusterConfig", + p: &priorityConfig{ + childNameGen: newNameGenerator(0), + clusterConfig: &xdsresource.ClusterConfig{}, + }, + }, + { + name: "EDS with nil EndpointConfig", + p: &priorityConfig{ + childNameGen: newNameGenerator(0), + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeEDS}, + }, + }, + }, + { + name: "EDS with nil EDSUpdate", + p: &priorityConfig{ + childNameGen: newNameGenerator(0), + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeEDS}, + EndpointConfig: &xdsresource.EndpointConfig{}, + }, + }, + }, + { + name: "LogicalDNS with nil EndpointConfig", + p: &priorityConfig{ + childNameGen: newNameGenerator(0), + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeLogicalDNS}, + }, + }, + }, + { + name: "LogicalDNS with nil DNSEndpoints", + p: &priorityConfig{ + childNameGen: newNameGenerator(0), + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeLogicalDNS}, + EndpointConfig: &xdsresource.EndpointConfig{}, + }, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, _, err := buildSingleClusterConfig(tt.p, nil) + if err == nil { + t.Errorf("buildSingleClusterConfig() succeeded unexpectedly, want error") + } + }) + } +} + func (s) TestBuildClusterImplConfigForDNS(t *testing.T) { for _, tt := range []struct { name string From c1ab2d9f3ac72360ed974074e8dd30b8c52214f2 Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Mon, 3 Aug 2026 12:26:01 +0000 Subject: [PATCH 04/10] Revert "Address Gemini review - fix nil checks" This reverts commit fbed93a2c516f848e7e45ec395d3a2e94236dced. --- .../xds/balancer/cdsbalancer/cdsbalancer.go | 12 +-- .../xds/balancer/cdsbalancer/configbuilder.go | 9 -- .../cdsbalancer/configbuilder_test.go | 82 ------------------- 3 files changed, 1 insertion(+), 102 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index d84df15eac20..6bcdd308bcda 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -273,14 +273,7 @@ func (b *cdsBalancer) handleClusterUpdate() error { // configuration is then pushed to the child policy. func (b *cdsBalancer) updateChildConfig() error { clusterName := b.lbCfg.ClusterName - clusterResult := b.clusterConfigs[clusterName] - if clusterResult == nil { - return fmt.Errorf("cluster result for %q is missing", clusterName) - } - clusterConfig := clusterResult.Config - if clusterConfig.Cluster == nil { - return fmt.Errorf("cluster update is missing for cluster %q", clusterName) - } + clusterConfig := b.clusterConfigs[clusterName].Config isAggregate := clusterConfig.Cluster.ClusterType == xdsresource.ClusterTypeAggregate var topLBName string @@ -312,9 +305,6 @@ func (b *cdsBalancer) updateChildConfig() error { if isAggregate { childCfgBytes, endpoints, err = buildAggregateClusterPriorityConfigJSON(b.priorities, &b.xdsLBPolicy) } else { - if len(b.priorities) == 0 { - return fmt.Errorf("no priorities configured for non-aggregate cluster %q", clusterName) - } childCfgBytes, endpoints, err = buildSingleClusterConfigJSON(b.priorities[0], &b.xdsLBPolicy) } if err != nil { diff --git a/internal/xds/balancer/cdsbalancer/configbuilder.go b/internal/xds/balancer/cdsbalancer/configbuilder.go index 322203390c24..4babc133645d 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder.go @@ -110,9 +110,6 @@ func buildSingleClusterConfigJSON(p *priorityConfig, xdsLBPolicy *internalservic } func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*outlierdetection.LBConfig, []resolver.Endpoint, error) { - if p == nil || p.clusterConfig == nil || p.clusterConfig.Cluster == nil || p.childNameGen == nil { - return nil, nil, fmt.Errorf("cluster config is missing") - } clusterUpdate := p.clusterConfig.Cluster priorityLBConfig := &priority.LBConfig{ Children: make(map[string]*priority.Child), @@ -122,9 +119,6 @@ func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalservicecon switch clusterUpdate.ClusterType { case xdsresource.ClusterTypeEDS: priorities := [][]xdsresource.Locality{{}} - if p.clusterConfig.EndpointConfig == nil || p.clusterConfig.EndpointConfig.EDSUpdate == nil { - return nil, nil, fmt.Errorf("EDS update is missing for cluster %q", clusterUpdate.ClusterName) - } edsUpdate := p.clusterConfig.EndpointConfig.EDSUpdate if len(edsUpdate.Localities) != 0 { priorities = groupLocalitiesByPriority(edsUpdate.Localities) @@ -147,9 +141,6 @@ func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalservicecon case xdsresource.ClusterTypeLogicalDNS: pName := fmt.Sprintf("priority-%v", p.childNameGen.prefix) priorityLBConfig.Priorities = []string{pName} - if p.clusterConfig.EndpointConfig == nil || p.clusterConfig.EndpointConfig.DNSEndpoints == nil { - return nil, nil, fmt.Errorf("DNS endpoints are missing for cluster %q", clusterUpdate.ClusterName) - } endpoints := p.clusterConfig.EndpointConfig.DNSEndpoints.Endpoints var retEndpoint resolver.Endpoint for _, e := range endpoints { diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 431e610480ba..7e93f940be16 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -681,88 +681,6 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t * } } -func (s) TestBuildSingleClusterConfig_NilChecks(t *testing.T) { - tests := []struct { - name string - p *priorityConfig - }{ - { - name: "nil priorityConfig", - p: nil, - }, - { - name: "priorityConfig with nil fields", - p: &priorityConfig{}, - }, - { - name: "nil clusterConfig", - p: &priorityConfig{childNameGen: newNameGenerator(0)}, - }, - { - name: "nil childNameGen", - p: &priorityConfig{ - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeEDS}, - }, - }, - }, - { - name: "nil Cluster in clusterConfig", - p: &priorityConfig{ - childNameGen: newNameGenerator(0), - clusterConfig: &xdsresource.ClusterConfig{}, - }, - }, - { - name: "EDS with nil EndpointConfig", - p: &priorityConfig{ - childNameGen: newNameGenerator(0), - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeEDS}, - }, - }, - }, - { - name: "EDS with nil EDSUpdate", - p: &priorityConfig{ - childNameGen: newNameGenerator(0), - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeEDS}, - EndpointConfig: &xdsresource.EndpointConfig{}, - }, - }, - }, - { - name: "LogicalDNS with nil EndpointConfig", - p: &priorityConfig{ - childNameGen: newNameGenerator(0), - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeLogicalDNS}, - }, - }, - }, - { - name: "LogicalDNS with nil DNSEndpoints", - p: &priorityConfig{ - childNameGen: newNameGenerator(0), - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ClusterType: xdsresource.ClusterTypeLogicalDNS}, - EndpointConfig: &xdsresource.EndpointConfig{}, - }, - }, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - _, _, err := buildSingleClusterConfig(tt.p, nil) - if err == nil { - t.Errorf("buildSingleClusterConfig() succeeded unexpectedly, want error") - } - }) - } -} - func (s) TestBuildClusterImplConfigForDNS(t *testing.T) { for _, tt := range []struct { name string From 785be23a1261c899177928f30912d1c9bdc55f6b Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Tue, 4 Aug 2026 08:30:58 +0000 Subject: [PATCH 05/10] Fix comments and legacy code cleanup --- .../xds/balancer/cdsbalancer/cdsbalancer.go | 28 ++++--------------- 1 file changed, 6 insertions(+), 22 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 6bcdd308bcda..960d4837806d 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -25,7 +25,6 @@ import ( "google.golang.org/grpc/balancer" "google.golang.org/grpc/balancer/base" "google.golang.org/grpc/connectivity" - "google.golang.org/grpc/internal/balancer/nop" "google.golang.org/grpc/internal/grpclog" "google.golang.org/grpc/internal/pretty" internalserviceconfig "google.golang.org/grpc/internal/serviceconfig" @@ -41,8 +40,8 @@ import ( const cdsName = "cds_experimental" var ( - // newChildBalancer is a helper function to build a new child balancer and its - // config parser, and will be overridden in unittests. + // newChildBalancer is a helper function to build a new child balancer + // and its config parser, and will be overridden in unittests. newChildBalancer = func(name string, cc balancer.ClientConn, opts balancer.BuildOptions) (balancer.Balancer, balancer.ConfigParser, error) { builder := balancer.Get(name) if builder == nil { @@ -69,26 +68,11 @@ type bb struct{} // Build creates a new CDS balancer with the ClientConn. func (bb) Build(cc balancer.ClientConn, opts balancer.BuildOptions) balancer.Balancer { - builder := balancer.Get(priority.Name) - if builder == nil { - // Shouldn't happen, registered through imported Priority builder. Still, - // defensive programming. - logger.Errorf("%q LB policy is needed but not registered", priority.Name) - return nop.NewBalancer(cc, fmt.Errorf("%q LB policy is needed but not registered", priority.Name)) - } - parser, ok := builder.(balancer.ConfigParser) - if !ok { - // Shouldn't happen, imported Priority builder has this method. - logger.Errorf("%q LB policy does not implement a config parser", priority.Name) - return nop.NewBalancer(cc, fmt.Errorf("%q LB policy does not implement a config parser", priority.Name)) - } - b := &cdsBalancer{ - bOpts: opts, - childConfigParser: parser, - clusterConfigs: make(map[string]*xdsresource.ClusterResult), - priorityConfigs: make(map[string]*priorityConfig), - cc: cc, + bOpts: opts, + clusterConfigs: make(map[string]*xdsresource.ClusterResult), + priorityConfigs: make(map[string]*priorityConfig), + cc: cc, } b.logger = prefixLogger(b) b.logger.Infof("Created") From 51e986ece5f8d878d1feed588a766573935319a1 Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Tue, 4 Aug 2026 08:54:55 +0000 Subject: [PATCH 06/10] Use gracefulswitch.Balancer --- .../xds/balancer/cdsbalancer/cdsbalancer.go | 52 +++++-------------- 1 file changed, 14 insertions(+), 38 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 960d4837806d..268c6ce87175 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -25,6 +25,7 @@ import ( "google.golang.org/grpc/balancer" "google.golang.org/grpc/balancer/base" "google.golang.org/grpc/connectivity" + "google.golang.org/grpc/internal/balancer/gracefulswitch" "google.golang.org/grpc/internal/grpclog" "google.golang.org/grpc/internal/pretty" internalserviceconfig "google.golang.org/grpc/internal/serviceconfig" @@ -39,24 +40,6 @@ import ( const cdsName = "cds_experimental" -var ( - // newChildBalancer is a helper function to build a new child balancer - // and its config parser, and will be overridden in unittests. - newChildBalancer = func(name string, cc balancer.ClientConn, opts balancer.BuildOptions) (balancer.Balancer, balancer.ConfigParser, error) { - builder := balancer.Get(name) - if builder == nil { - return nil, nil, fmt.Errorf("xds: no balancer builder with name %v", name) - } - parser, ok := builder.(balancer.ConfigParser) - if !ok { - return nil, nil, fmt.Errorf("xds: balancer builder for %v does not implement ConfigParser", name) - } - // We directly pass the parent clientConn to the underlying child - // balancer because the cdsBalancer does not deal with subConns. - return builder.Build(cc, opts), parser, nil - } -) - func init() { balancer.Register(bb{}) } @@ -70,6 +53,7 @@ type bb struct{} func (bb) Build(cc balancer.ClientConn, opts balancer.BuildOptions) balancer.Balancer { b := &cdsBalancer{ bOpts: opts, + childLB: gracefulswitch.NewBalancer(cc, opts), clusterConfigs: make(map[string]*xdsresource.ClusterResult), priorityConfigs: make(map[string]*priorityConfig), cc: cc, @@ -111,18 +95,16 @@ type cdsBalancer struct { // The following fields are initialized at build time and are either // read-only after that or provide their own synchronization, and therefore // do not need to be guarded by a mutex. - cc balancer.ClientConn // ClientConn interface passed to child LB. - bOpts balancer.BuildOptions // BuildOptions passed to child LB. - childConfigParser balancer.ConfigParser // Config parser for cluster_resolver LB policy. - logger *grpclog.PrefixLogger // Prefix logger for all logging. + cc balancer.ClientConn // ClientConn interface passed to child LB. + bOpts balancer.BuildOptions // BuildOptions passed to child LB. + logger *grpclog.PrefixLogger // Prefix logger for all logging. // All fields below are accessed only from methods implementing the // balancer.Balancer interface. Since gRPC guarantees that these methods are // never invoked concurrently, no additional synchronization is required to // protect access to these fields. xdsClient xdsclient.XDSClient - childLB balancer.Balancer // Child policy, built upon resolution of the cluster graph. - childLBName string // Name of the child policy. + childLB *gracefulswitch.Balancer // Graceful switch child policy. clusterConfigs map[string]*xdsresource.ClusterResult // Cluster name to the last received result for that cluster. priorityConfigs map[string]*priorityConfig // Hostname to priority config for that leaf cluster. lbCfg *lbConfig // Current load balancing configuration. @@ -267,19 +249,8 @@ func (b *cdsBalancer) updateChildConfig() error { topLBName = outlierdetection.Name } - if b.childLB != nil && b.childLBName != topLBName { - b.childLB.Close() - b.childLB = nil - } - if b.childLB == nil { - childLB, parser, err := newChildBalancer(topLBName, b.cc, b.bOpts) - if err != nil { - return fmt.Errorf("failed to create child policy of type %s: %v", topLBName, err) - } - b.childLB = childLB - b.childLBName = topLBName - b.childConfigParser = parser + b.childLB = gracefulswitch.NewBalancer(b.cc, b.bOpts) } var childCfgBytes []byte @@ -295,9 +266,14 @@ func (b *cdsBalancer) updateChildConfig() error { return fmt.Errorf("failed to build child policy config: %v", err) } - childCfg, err := b.childConfigParser.ParseConfig(childCfgBytes) + cfgJSON, err := json.Marshal([]map[string]json.RawMessage{{topLBName: childCfgBytes}}) + if err != nil { + return fmt.Errorf("failed to marshal child policy config wrapper: %v", err) + } + + childCfg, err := gracefulswitch.ParseConfig(cfgJSON) if err != nil { - return fmt.Errorf("failed to parse child policy config. This should never happen because the config was generated: %v", err) + return fmt.Errorf("failed to parse child policy config: %v", err) } if b.logger.V(2) { b.logger.Infof("Built child policy config: %s", pretty.ToJSON(childCfg)) From e7c04757e5718cfa53d14f662bf82dae1f375587 Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Tue, 4 Aug 2026 09:13:21 +0000 Subject: [PATCH 07/10] Fix functions and variables names --- .../cdsbalancer/aggregate_cluster_test.go | 18 +++---- .../xds/balancer/cdsbalancer/cdsbalancer.go | 4 +- .../balancer/cdsbalancer/cdsbalancer_test.go | 6 +-- .../xds/balancer/cdsbalancer/configbuilder.go | 18 +++---- .../cdsbalancer/configbuilder_test.go | 52 +++++++++---------- 5 files changed, 49 insertions(+), 49 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go b/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go index b56bab30186d..595673bd6a0c 100644 --- a/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go +++ b/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go @@ -105,15 +105,15 @@ func (s) TestAggregateClusterSuccess_LeafNode(t *testing.T) { name: "eds", firstClusterResource: e2e.DefaultCluster(clusterName, serviceName, e2e.SecurityLevelNone), secondClusterResource: e2e.DefaultCluster(clusterName, serviceName+"-new", e2e.SecurityLevelNone), - wantFirstChildCfg: createSingleClusterConfig(clusterName, "priority-0-0", true), - wantSecondChildCfg: createSingleClusterConfig(clusterName, "priority-1-0", true), + wantFirstChildCfg: createLeafClusterConfig(clusterName, "priority-0-0", true), + wantSecondChildCfg: createLeafClusterConfig(clusterName, "priority-1-0", true), }, { name: "dns", firstClusterResource: makeLogicalDNSClusterResource(clusterName, "dns_host", uint32(port)), secondClusterResource: makeLogicalDNSClusterResource(clusterName, "dns_host_new", uint32(port)), - wantFirstChildCfg: createSingleClusterConfig(clusterName, "priority-0", false), - wantSecondChildCfg: createSingleClusterConfig(clusterName, "priority-1", false), + wantFirstChildCfg: createLeafClusterConfig(clusterName, "priority-0", false), + wantSecondChildCfg: createLeafClusterConfig(clusterName, "priority-1", false), }, } @@ -319,8 +319,8 @@ func (s) TestAggregateClusterSuccess_ThenChangeRootToEDS(t *testing.T) { } // Since the service name of the EDS cluster remains same, same priority name // is used. - wantSingleChildCfg := createSingleClusterConfig(clusterName, "priority-0-0", true) - if err := compareLoadBalancingConfig(ctx, odCfgCh, wantSingleChildCfg); err != nil { + wantLeafChildCfg := createLeafClusterConfig(clusterName, "priority-0-0", true) + if err := compareLoadBalancingConfig(ctx, odCfgCh, wantLeafChildCfg); err != nil { t.Fatal(err) } } @@ -348,8 +348,8 @@ func (s) TestAggregatedClusterSuccess_SwitchBetweenLeafAndAggregate(t *testing.T if err := mgmtServer.Update(ctx, resources); err != nil { t.Fatal(err) } - wantSingleChildCfg := createSingleClusterConfig(clusterName, "priority-0-0", true) - if err := compareLoadBalancingConfig(ctx, odCfgCh, wantSingleChildCfg); err != nil { + wantLeafChildCfg := createLeafClusterConfig(clusterName, "priority-0-0", true) + if err := compareLoadBalancingConfig(ctx, odCfgCh, wantLeafChildCfg); err != nil { t.Fatal(err) } @@ -400,7 +400,7 @@ func (s) TestAggregatedClusterSuccess_SwitchBetweenLeafAndAggregate(t *testing.T if err := mgmtServer.Update(ctx, resources); err != nil { t.Fatal(err) } - if err := compareLoadBalancingConfig(ctx, odCfgCh, wantSingleChildCfg); err != nil { + if err := compareLoadBalancingConfig(ctx, odCfgCh, wantLeafChildCfg); err != nil { t.Fatal(err) } } diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 268c6ce87175..62d515b47438 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -258,9 +258,9 @@ func (b *cdsBalancer) updateChildConfig() error { var err error if isAggregate { - childCfgBytes, endpoints, err = buildAggregateClusterPriorityConfigJSON(b.priorities, &b.xdsLBPolicy) + childCfgBytes, endpoints, err = buildAggregateClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } else { - childCfgBytes, endpoints, err = buildSingleClusterConfigJSON(b.priorities[0], &b.xdsLBPolicy) + childCfgBytes, endpoints, err = buildLeafClusterConfigJSON(b.priorities[0], &b.xdsLBPolicy) } if err != nil { return fmt.Errorf("failed to build child policy config: %v", err) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go b/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go index 77a787c68aa9..38ab28132053 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer_test.go @@ -337,10 +337,10 @@ func createPriorityConfig(cluster string) *iserviceconfig.BalancerConfig { } } -// createSingleClusterConfig returns the expected LoadBalancingConfig tree for a -// single (non-aggregate) cluster under gRFC A75 topology: +// createLeafClusterConfig returns the expected LoadBalancingConfig tree for a +// leaf (non-aggregate) cluster under gRFC A75 topology: // outlier_detection -> cluster_impl -> priority -> wrr_locality -> round_robin. -func createSingleClusterConfig(cluster string, pName string, ignoreReresolution bool) *outlierdetection.LBConfig { +func createLeafClusterConfig(cluster string, pName string, ignoreReresolution bool) *outlierdetection.LBConfig { return &outlierdetection.LBConfig{ Interval: iserviceconfig.Duration(10 * time.Second), BaseEjectionTime: iserviceconfig.Duration(30 * time.Second), diff --git a/internal/xds/balancer/cdsbalancer/configbuilder.go b/internal/xds/balancer/cdsbalancer/configbuilder.go index 4babc133645d..d63f6cd7061c 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder.go @@ -77,7 +77,7 @@ func hostName(clusterName string, update xdsresource.ClusterUpdate) string { } } -// buildSingleClusterConfigJSON builds the balancer config for a single +// buildLeafClusterConfigJSON builds the balancer config for a leaf // (non-aggregate) cluster according to gRFC A75. // // The built tree of balancers: @@ -97,19 +97,19 @@ func hostName(clusterName string, update xdsresource.ClusterUpdate) string { // ┌──────────▼─┐ ┌─▼──────────┐ // │xDSLBPolicy │ │xDSLBPolicy │ (Locality and Endpoint picking layer) // └────────────┘ └────────────┘ -func buildSingleClusterConfigJSON(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { - odCfg, endpoints, err := buildSingleClusterConfig(p, xdsLBPolicy) +func buildLeafClusterConfigJSON(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { + odCfg, endpoints, err := buildLeafClusterConfig(p, xdsLBPolicy) if err != nil { return nil, nil, err } ret, err := json.Marshal(odCfg) if err != nil { - return nil, nil, fmt.Errorf("failed to marshal built single cluster config: %v", err) + return nil, nil, fmt.Errorf("failed to marshal built leaf cluster config: %v", err) } return ret, endpoints, nil } -func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*outlierdetection.LBConfig, []resolver.Endpoint, error) { +func buildLeafClusterConfig(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*outlierdetection.LBConfig, []resolver.Endpoint, error) { clusterUpdate := p.clusterConfig.Cluster priorityLBConfig := &priority.LBConfig{ Children: make(map[string]*priority.Child), @@ -173,7 +173,7 @@ func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalservicecon return &odCfg, retEndpoints, nil } -// buildAggregateClusterPriorityConfigJSON builds balancer config for the passed in +// buildAggregateClusterConfigJSON builds balancer config for the passed in // priorities (legacy / aggregate cluster tree). // // The built tree of balancers: @@ -189,8 +189,8 @@ func buildSingleClusterConfig(p *priorityConfig, xdsLBPolicy *internalservicecon // ┌──────▼─────┐ ┌─────▼──────┐ // │xDSLBPolicy │ │xDSLBPolicy │ (Locality and Endpoint picking layer) // └────────────┘ └────────────┘ -func buildAggregateClusterPriorityConfigJSON(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { - pc, endpoints, err := buildAggregateClusterPriorityConfig(priorities, xdsLBPolicy) +func buildAggregateClusterConfigJSON(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { + pc, endpoints, err := buildAggregateClusterConfig(priorities, xdsLBPolicy) if err != nil { return nil, nil, fmt.Errorf("failed to build priority config: %v", err) } @@ -201,7 +201,7 @@ func buildAggregateClusterPriorityConfigJSON(priorities []*priorityConfig, xdsLB return ret, endpoints, nil } -func buildAggregateClusterPriorityConfig(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*priority.LBConfig, []resolver.Endpoint, error) { +func buildAggregateClusterConfig(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*priority.LBConfig, []resolver.Endpoint, error) { var ( retConfig = &priority.LBConfig{Children: make(map[string]*priority.Child)} retEndpoints []resolver.Endpoint diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 7e93f940be16..05790d057c77 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -130,9 +130,9 @@ func makeLocality(localityIdx int, localityWeight, priority uint32, endpointCoun } } -// TestBuildAggregateClusterPriorityConfigJSON is a sanity check that the built balancer config -// can be parsed. The behavior test is covered by TestBuildAggregateClusterPriorityConfig. -func (s) TestBuildAggregateClusterPriorityConfigJSON(t *testing.T) { +// TestBuildAggregateClusterConfigJSON is a sanity check that the built balancer config +// can be parsed. The behavior test is covered by TestBuildAggregateClusterConfig. +func (s) TestBuildAggregateClusterConfigJSON(t *testing.T) { testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ URI: "trafficdirector.googleapis.com:443", ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, @@ -141,7 +141,7 @@ func (s) TestBuildAggregateClusterPriorityConfigJSON(t *testing.T) { t.Fatalf("Failed to create LRS server config for testing: %v", err) } - gotConfig, _, err := buildAggregateClusterPriorityConfigJSON([]*priorityConfig{ + gotConfig, _, err := buildAggregateClusterConfigJSON([]*priorityConfig{ { clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ @@ -182,7 +182,7 @@ func (s) TestBuildAggregateClusterPriorityConfigJSON(t *testing.T) { }, }, nil) if err != nil { - t.Fatalf("buildAggregateClusterPriorityConfigJSON(...) failed: %v", err) + t.Fatalf("buildAggregateClusterConfigJSON(...) failed: %v", err) } var prettyGot bytes.Buffer @@ -198,11 +198,11 @@ func (s) TestBuildAggregateClusterPriorityConfigJSON(t *testing.T) { } } -// TestBuildAggregateClusterPriorityConfig tests the priority config generation. Each top level +// TestBuildAggregateClusterConfig tests the priority config generation. Each top level // balancer per priority should be an Outlier Detection balancer, with a Cluster // Impl Balancer as a child. -func (s) TestBuildAggregateClusterPriorityConfig(t *testing.T) { - gotConfig, _, _ := buildAggregateClusterPriorityConfig([]*priorityConfig{ +func (s) TestBuildAggregateClusterConfig(t *testing.T) { + gotConfig, _, _ := buildAggregateClusterConfig([]*priorityConfig{ { // EDS - OD config should be the top level for both of the EDS // priorities balancer This EDS priority will have multiple sub @@ -317,7 +317,7 @@ func testEndpointForDNS(endpoints []resolver.Endpoint, localityWeight uint32, pa return retEndpoint } -func (s) TestBuildSingleClusterConfigJSON(t *testing.T) { +func (s) TestBuildLeafClusterConfigJSON(t *testing.T) { testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ URI: "trafficdirector.googleapis.com:443", ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, @@ -326,7 +326,7 @@ func (s) TestBuildSingleClusterConfigJSON(t *testing.T) { t.Fatalf("Failed to create LRS server config for testing: %v", err) } - gotConfig, _, err := buildSingleClusterConfigJSON(&priorityConfig{ + gotConfig, _, err := buildLeafClusterConfigJSON(&priorityConfig{ clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ ClusterName: testClusterName, @@ -355,7 +355,7 @@ func (s) TestBuildSingleClusterConfigJSON(t *testing.T) { childNameGen: newNameGenerator(0), }, &iserviceconfig.BalancerConfig{Name: roundrobin.Name}) if err != nil { - t.Fatalf("buildSingleClusterConfigJSON(...) failed: %v", err) + t.Fatalf("buildLeafClusterConfigJSON(...) failed: %v", err) } var prettyGot bytes.Buffer @@ -370,7 +370,7 @@ func (s) TestBuildSingleClusterConfigJSON(t *testing.T) { } } -func (s) TestBuildSingleClusterConfig_DNS(t *testing.T) { +func (s) TestBuildLeafClusterConfig_DNS(t *testing.T) { for _, tt := range []struct { name string endpoints []resolver.Endpoint @@ -413,7 +413,7 @@ func (s) TestBuildSingleClusterConfig_DNS(t *testing.T) { }, } { t.Run(tt.name, func(t *testing.T) { - gotODConfig, gotEndpoints, err := buildSingleClusterConfig( + gotODConfig, gotEndpoints, err := buildLeafClusterConfig( &priorityConfig{ clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ @@ -427,7 +427,7 @@ func (s) TestBuildSingleClusterConfig_DNS(t *testing.T) { }, tt.xdsLBPolicy) if err != nil { - t.Fatalf("buildSingleClusterConfig() failed: %v", err) + t.Fatalf("buildLeafClusterConfig() failed: %v", err) } wantODConfig := &outlierdetection.LBConfig{ @@ -455,18 +455,18 @@ func (s) TestBuildSingleClusterConfig_DNS(t *testing.T) { }, } if diff := cmp.Diff(wantODConfig, gotODConfig); diff != "" { - t.Errorf("buildSingleClusterConfig() config diff (-want +got) %v", diff) + t.Errorf("buildLeafClusterConfig() config diff (-want +got) %v", diff) } wantEndpoints := []resolver.Endpoint{testEndpointForDNS(tt.endpoints, 1, []string{"priority-3", xdsinternal.LocalityString(clients.Locality{})})} if diff := cmp.Diff(wantEndpoints, gotEndpoints, endpointCmpOpts); diff != "" { - t.Errorf("buildSingleClusterConfig() endpoints diff (-want +got) %v", diff) + t.Errorf("buildLeafClusterConfig() endpoints diff (-want +got) %v", diff) } }) } } -func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t *testing.T) { +func (s) TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, false) testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ @@ -484,7 +484,7 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t loc2 := makeLocality(2, 20, 1, 2) loc3 := makeLocality(3, 80, 1, 2) - gotODConfig, gotEndpoints, err := buildSingleClusterConfig( + gotODConfig, gotEndpoints, err := buildLeafClusterConfig( &priorityConfig{ clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ @@ -516,7 +516,7 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t nil, ) if err != nil { - t.Fatalf("buildSingleClusterConfig() failed: %v", err) + t.Fatalf("buildLeafClusterConfig() failed: %v", err) } wantODConfig := &outlierdetection.LBConfig{ @@ -558,14 +558,14 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t } if diff := cmp.Diff(wantODConfig, gotODConfig); diff != "" { - t.Errorf("buildSingleClusterConfig() config diff (-want +got) %v", diff) + t.Errorf("buildLeafClusterConfig() config diff (-want +got) %v", diff) } if diff := cmp.Diff(wantEndpoints, gotEndpoints, endpointCmpOpts); diff != "" { - t.Errorf("buildSingleClusterConfig() endpoints diff (-want +got) %v", diff) + t.Errorf("buildLeafClusterConfig() endpoints diff (-want +got) %v", diff) } } -func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t *testing.T) { +func (s) TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, true) testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ @@ -583,7 +583,7 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t * loc2 := makeLocality(2, 20, 1, 2) loc3 := makeLocality(3, 80, 1, 2) - gotODConfig, gotEndpoints, err := buildSingleClusterConfig( + gotODConfig, gotEndpoints, err := buildLeafClusterConfig( &priorityConfig{ clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ @@ -615,7 +615,7 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t * nil, ) if err != nil { - t.Fatalf("buildSingleClusterConfig() failed: %v", err) + t.Fatalf("buildLeafClusterConfig() failed: %v", err) } wantODConfig := &outlierdetection.LBConfig{ @@ -674,10 +674,10 @@ func (s) TestBuildSingleClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t * } if diff := cmp.Diff(wantODConfig, gotODConfig); diff != "" { - t.Errorf("buildSingleClusterConfig() config diff (-want +got) %v", diff) + t.Errorf("buildLeafClusterConfig() config diff (-want +got) %v", diff) } if diff := cmp.Diff(wantEndpoints, gotEndpoints, endpointCmpOpts); diff != "" { - t.Errorf("buildSingleClusterConfig() endpoints diff (-want +got) %v", diff) + t.Errorf("buildLeafClusterConfig() endpoints diff (-want +got) %v", diff) } } From d6a83ca10ca94f4e8792b3e328b871d2263169de Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Tue, 4 Aug 2026 10:41:57 +0000 Subject: [PATCH 08/10] Do not create clusterImpl for each priority --- .../xds/balancer/cdsbalancer/cdsbalancer.go | 2 +- .../xds/balancer/cdsbalancer/configbuilder.go | 26 ++++++++++++------- .../cdsbalancer/configbuilder_test.go | 4 +-- 3 files changed, 19 insertions(+), 13 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 62d515b47438..9637f7f752f1 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -260,7 +260,7 @@ func (b *cdsBalancer) updateChildConfig() error { if isAggregate { childCfgBytes, endpoints, err = buildAggregateClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } else { - childCfgBytes, endpoints, err = buildLeafClusterConfigJSON(b.priorities[0], &b.xdsLBPolicy) + childCfgBytes, endpoints, err = buildLeafClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } if err != nil { return fmt.Errorf("failed to build child policy config: %v", err) diff --git a/internal/xds/balancer/cdsbalancer/configbuilder.go b/internal/xds/balancer/cdsbalancer/configbuilder.go index d63f6cd7061c..65e610b93766 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder.go @@ -97,8 +97,8 @@ func hostName(clusterName string, update xdsresource.ClusterUpdate) string { // ┌──────────▼─┐ ┌─▼──────────┐ // │xDSLBPolicy │ │xDSLBPolicy │ (Locality and Endpoint picking layer) // └────────────┘ └────────────┘ -func buildLeafClusterConfigJSON(p *priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { - odCfg, endpoints, err := buildLeafClusterConfig(p, xdsLBPolicy) +func buildLeafClusterConfigJSON(priorities []*priorityConfig, xdsLBPolicy *internalserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) { + odCfg, endpoints, err := buildLeafClusterConfig(priorities[0], xdsLBPolicy) if err != nil { return nil, nil, err } @@ -128,10 +128,7 @@ func buildLeafClusterConfig(p *priorityConfig, xdsLBPolicy *internalserviceconfi for i, pName := range priorityNames { priorityLocalities := priorities[i] - _, endpoints, err := priorityLocalitiesToClusterImpl(priorityLocalities, pName, *p.clusterConfig.Cluster, xdsLBPolicy) - if err != nil { - return nil, nil, err - } + endpoints := priorityLocalitiesToEndpoints(priorityLocalities, pName, *p.clusterConfig.Cluster) retEndpoints = append(retEndpoints, endpoints...) priorityLBConfig.Children[pName] = &priority.Child{ Config: xdsLBPolicy, @@ -354,6 +351,18 @@ func groupLocalitiesByPriority(localities []xdsresource.Locality) [][]xdsresourc // addresses with their path hierarchy set to [priority-name, locality-name], so // priority and the xDS LB Policy know which child policy each address is for. func priorityLocalitiesToClusterImpl(localities []xdsresource.Locality, priorityName string, clusterUpdate xdsresource.ClusterUpdate, xdsLBPolicy *internalserviceconfig.BalancerConfig) (*clusterimpl.LBConfig, []resolver.Endpoint, error) { + endpoints := priorityLocalitiesToEndpoints(localities, priorityName, clusterUpdate) + return &clusterimpl.LBConfig{ + Cluster: clusterUpdate.ClusterName, + ChildPolicy: xdsLBPolicy, + }, endpoints, nil +} + +// priorityLocalitiesToEndpoints takes a list of localities (with the same +// priority), and generates a list of addresses with their path hierarchy set to +// [priority-name, locality-name], so priority and the xDS LB Policy know which +// child policy each address is for. +func priorityLocalitiesToEndpoints(localities []xdsresource.Locality, priorityName string, clusterUpdate xdsresource.ClusterUpdate) []resolver.Endpoint { var retEndpoints []resolver.Endpoint // Compute the sum of locality weights to normalize locality weights. The @@ -432,10 +441,7 @@ func priorityLocalitiesToClusterImpl(localities []xdsresource.Locality, priority retEndpoints = append(retEndpoints, resolverEndpoint) } } - return &clusterimpl.LBConfig{ - Cluster: clusterUpdate.ClusterName, - ChildPolicy: xdsLBPolicy, - }, retEndpoints, nil + return retEndpoints } // fixedPointFractionalBits is the number of bits used for the fractional part diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 05790d057c77..8c58b04aad7d 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -326,7 +326,7 @@ func (s) TestBuildLeafClusterConfigJSON(t *testing.T) { t.Fatalf("Failed to create LRS server config for testing: %v", err) } - gotConfig, _, err := buildLeafClusterConfigJSON(&priorityConfig{ + gotConfig, _, err := buildLeafClusterConfigJSON([]*priorityConfig{{ clusterConfig: &xdsresource.ClusterConfig{ Cluster: &xdsresource.ClusterUpdate{ ClusterName: testClusterName, @@ -353,7 +353,7 @@ func (s) TestBuildLeafClusterConfigJSON(t *testing.T) { }, outlierDetection: noopODCfg, childNameGen: newNameGenerator(0), - }, &iserviceconfig.BalancerConfig{Name: roundrobin.Name}) + }}, &iserviceconfig.BalancerConfig{Name: roundrobin.Name}) if err != nil { t.Fatalf("buildLeafClusterConfigJSON(...) failed: %v", err) } From 9f43e37a99edb8e277db4b2d6a0730afaffd787f Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Wed, 12 Aug 2026 11:37:34 +0000 Subject: [PATCH 09/10] Address feedback --- .../cdsbalancer/aggregate_cluster_test.go | 11 +- .../xds/balancer/cdsbalancer/cdsbalancer.go | 10 +- .../cdsbalancer/configbuilder_test.go | 222 ++++++++++-------- 3 files changed, 134 insertions(+), 109 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go b/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go index 595673bd6a0c..f565ed8c56a7 100644 --- a/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go +++ b/internal/xds/balancer/cdsbalancer/aggregate_cluster_test.go @@ -90,8 +90,8 @@ func verifyDNSResolution(ctx context.Context, t *testing.T, dnsTargetCh chan res // Tests the case where the cluster resource requested is a leaf cluster. The // management server sends two updates for the same leaf cluster resource. The -// test verifies that the load balancing configuration pushed to the top-level LB -// policy contains the expected configuration corresponding to the leaf +// test verifies that the load balancing configuration pushed to the top-level +// LB policy contains the expected configuration corresponding to the leaf // cluster, on both occasions. func (s) TestAggregateClusterSuccess_LeafNode(t *testing.T) { tests := []struct { @@ -255,11 +255,12 @@ func (s) TestAggregateClusterSuccess_ThenUpdateChildClusters(t *testing.T) { // Tests the case where the cluster resource requested is an aggregate cluster // root pointing to two child clusters, one of type EDS and the other of type -// LogicalDNS. The test verifies that the load balancing configuration pushed to -// the priority LB policy contains the discovery mechanisms for both child +// LogicalDNS. The test verifies that the load balancing configuration pushed +// to the priority LB policy contains the discovery mechanisms for both child // clusters. The test then updates the root cluster resource requested by the // cds LB policy to a leaf cluster of type EDS and verifies the load balancing -// configuration pushed to the outlier detection LB policy contains the leaf cluster config. +// configuration pushed to the outlier detection LB policy contains the leaf +// cluster config. func (s) TestAggregateClusterSuccess_ThenChangeRootToEDS(t *testing.T) { dnsTargetCh, dnsR := setupDNS(t) lbCfgCh, _, _, _ := registerWrappedPriorityPolicy(t) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 9637f7f752f1..1e155358898d 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -242,13 +242,6 @@ func (b *cdsBalancer) updateChildConfig() error { clusterConfig := b.clusterConfigs[clusterName].Config isAggregate := clusterConfig.Cluster.ClusterType == xdsresource.ClusterTypeAggregate - var topLBName string - if isAggregate { - topLBName = priority.Name - } else { - topLBName = outlierdetection.Name - } - if b.childLB == nil { b.childLB = gracefulswitch.NewBalancer(b.cc, b.bOpts) } @@ -256,10 +249,13 @@ func (b *cdsBalancer) updateChildConfig() error { var childCfgBytes []byte var endpoints []resolver.Endpoint var err error + var topLBName string if isAggregate { + topLBName = priority.Name childCfgBytes, endpoints, err = buildAggregateClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } else { + topLBName = outlierdetection.Name childCfgBytes, endpoints, err = buildLeafClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } if err != nil { diff --git a/internal/xds/balancer/cdsbalancer/configbuilder_test.go b/internal/xds/balancer/cdsbalancer/configbuilder_test.go index 8c58b04aad7d..e353fbccfbec 100644 --- a/internal/xds/balancer/cdsbalancer/configbuilder_test.go +++ b/internal/xds/balancer/cdsbalancer/configbuilder_test.go @@ -130,9 +130,10 @@ func makeLocality(localityIdx int, localityWeight, priority uint32, endpointCoun } } -// TestBuildAggregateClusterConfigJSON is a sanity check that the built balancer config -// can be parsed. The behavior test is covered by TestBuildAggregateClusterConfig. -func (s) TestBuildAggregateClusterConfigJSON(t *testing.T) { +// TestBuildClusterConfigJSON is a sanity check that built balancer configs +// for aggregate and leaf clusters can be parsed by their respective policy +// parsers. +func (s) TestBuildClusterConfigJSON(t *testing.T) { testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ URI: "trafficdirector.googleapis.com:443", ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, @@ -141,60 +142,113 @@ func (s) TestBuildAggregateClusterConfigJSON(t *testing.T) { t.Fatalf("Failed to create LRS server config for testing: %v", err) } - gotConfig, _, err := buildAggregateClusterConfigJSON([]*priorityConfig{ + tests := []struct { + name string + buildFunc func([]*priorityConfig, *iserviceconfig.BalancerConfig) ([]byte, []resolver.Endpoint, error) + priorities []*priorityConfig + xdsLBPolicy *iserviceconfig.BalancerConfig + policyName string + }{ { - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ - ClusterName: testClusterName, - ClusterType: xdsresource.ClusterTypeEDS, - EDSServiceName: testEDSServiceName, - MaxRequests: newUint32(testMaxRequests), - LRSServerConfig: testLRSServerConfig, + name: "aggregate", + buildFunc: buildAggregateClusterConfigJSON, + priorities: []*priorityConfig{ + { + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterName: testClusterName, + ClusterType: xdsresource.ClusterTypeEDS, + EDSServiceName: testEDSServiceName, + MaxRequests: newUint32(testMaxRequests), + LRSServerConfig: testLRSServerConfig, + }, + EndpointConfig: &xdsresource.EndpointConfig{ + EDSUpdate: &xdsresource.EndpointsUpdate{ + Drops: []xdsresource.OverloadDropConfig{{ + Category: testDropCategory, + Numerator: testDropOverMillion, + Denominator: million, + }}, + Localities: []xdsresource.Locality{ + makeLocality(0, 20, 0, 2), + makeLocality(1, 80, 0, 2), + makeLocality(2, 20, 1, 2), + makeLocality(3, 80, 1, 2), + }, + }, + }, + }, + childNameGen: newNameGenerator(0), }, - EndpointConfig: &xdsresource.EndpointConfig{ - EDSUpdate: &xdsresource.EndpointsUpdate{ - Drops: []xdsresource.OverloadDropConfig{{ - Category: testDropCategory, - Numerator: testDropOverMillion, - Denominator: million, - }}, - Localities: []xdsresource.Locality{ - makeLocality(0, 20, 0, 2), - makeLocality(1, 80, 0, 2), - makeLocality(2, 20, 1, 2), - makeLocality(3, 80, 1, 2), + { + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterType: xdsresource.ClusterTypeLogicalDNS, + }, + EndpointConfig: &xdsresource.EndpointConfig{ + DNSEndpoints: &xdsresource.DNSUpdate{Endpoints: []resolver.Endpoint{makeResolverEndpoint(4, 0), makeResolverEndpoint(4, 1)}}, }, }, + childNameGen: newNameGenerator(1), }, }, - childNameGen: newNameGenerator(0), + xdsLBPolicy: nil, + policyName: priority.Name, }, { - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ - ClusterType: xdsresource.ClusterTypeLogicalDNS, - }, - EndpointConfig: &xdsresource.EndpointConfig{ - DNSEndpoints: &xdsresource.DNSUpdate{Endpoints: []resolver.Endpoint{makeResolverEndpoint(4, 0), makeResolverEndpoint(4, 1)}}, + name: "leaf", + buildFunc: buildLeafClusterConfigJSON, + priorities: []*priorityConfig{{ + clusterConfig: &xdsresource.ClusterConfig{ + Cluster: &xdsresource.ClusterUpdate{ + ClusterName: testClusterName, + ClusterType: xdsresource.ClusterTypeEDS, + EDSServiceName: testEDSServiceName, + MaxRequests: newUint32(testMaxRequests), + LRSServerConfig: testLRSServerConfig, + }, + EndpointConfig: &xdsresource.EndpointConfig{ + EDSUpdate: &xdsresource.EndpointsUpdate{ + Drops: []xdsresource.OverloadDropConfig{{ + Category: testDropCategory, + Numerator: testDropOverMillion, + Denominator: million, + }}, + Localities: []xdsresource.Locality{ + makeLocality(0, 20, 0, 2), + makeLocality(1, 80, 0, 2), + makeLocality(2, 20, 1, 2), + makeLocality(3, 80, 1, 2), + }, + }, + }, }, - }, - childNameGen: newNameGenerator(1), + outlierDetection: noopODCfg, + childNameGen: newNameGenerator(0), + }}, + xdsLBPolicy: &iserviceconfig.BalancerConfig{Name: roundrobin.Name}, + policyName: outlierdetection.Name, }, - }, nil) - if err != nil { - t.Fatalf("buildAggregateClusterConfigJSON(...) failed: %v", err) } - var prettyGot bytes.Buffer - if err := json.Indent(&prettyGot, gotConfig, ">>> ", " "); err != nil { - t.Fatalf("json.Indent() failed: %v", err) - } - // Print the indented json if this test fails. - t.Log(prettyGot.String()) + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + gotConfig, _, err := tt.buildFunc(tt.priorities, tt.xdsLBPolicy) + if err != nil { + t.Fatalf("%s(...) failed: %v", tt.name, err) + } - priorityB := balancer.Get(priority.Name) - if _, err = priorityB.(balancer.ConfigParser).ParseConfig(gotConfig); err != nil { - t.Fatalf("ParseConfig(%+v) failed: %v", gotConfig, err) + var prettyGot bytes.Buffer + if err := json.Indent(&prettyGot, gotConfig, ">>> ", " "); err != nil { + t.Fatalf("json.Indent() failed: %v", err) + } + t.Log(prettyGot.String()) + + b := balancer.Get(tt.policyName) + if _, err = b.(balancer.ConfigParser).ParseConfig(gotConfig); err != nil { + t.Fatalf("ParseConfig(%+v) failed: %v", gotConfig, err) + } + }) } } @@ -317,59 +371,9 @@ func testEndpointForDNS(endpoints []resolver.Endpoint, localityWeight uint32, pa return retEndpoint } -func (s) TestBuildLeafClusterConfigJSON(t *testing.T) { - testLRSServerConfig, err := bootstrap.ServerConfigForTesting(bootstrap.ServerConfigTestingOptions{ - URI: "trafficdirector.googleapis.com:443", - ChannelCreds: []bootstrap.ChannelCreds{{Type: "google_default"}}, - }) - if err != nil { - t.Fatalf("Failed to create LRS server config for testing: %v", err) - } - - gotConfig, _, err := buildLeafClusterConfigJSON([]*priorityConfig{{ - clusterConfig: &xdsresource.ClusterConfig{ - Cluster: &xdsresource.ClusterUpdate{ - ClusterName: testClusterName, - ClusterType: xdsresource.ClusterTypeEDS, - EDSServiceName: testEDSServiceName, - MaxRequests: newUint32(testMaxRequests), - LRSServerConfig: testLRSServerConfig, - }, - EndpointConfig: &xdsresource.EndpointConfig{ - EDSUpdate: &xdsresource.EndpointsUpdate{ - Drops: []xdsresource.OverloadDropConfig{{ - Category: testDropCategory, - Numerator: testDropOverMillion, - Denominator: million, - }}, - Localities: []xdsresource.Locality{ - makeLocality(0, 20, 0, 2), - makeLocality(1, 80, 0, 2), - makeLocality(2, 20, 1, 2), - makeLocality(3, 80, 1, 2), - }, - }, - }, - }, - outlierDetection: noopODCfg, - childNameGen: newNameGenerator(0), - }}, &iserviceconfig.BalancerConfig{Name: roundrobin.Name}) - if err != nil { - t.Fatalf("buildLeafClusterConfigJSON(...) failed: %v", err) - } - - var prettyGot bytes.Buffer - if err := json.Indent(&prettyGot, gotConfig, ">>> ", " "); err != nil { - t.Fatalf("json.Indent() failed: %v", err) - } - t.Log(prettyGot.String()) - - odB := balancer.Get(outlierdetection.Name) - if _, err = odB.(balancer.ConfigParser).ParseConfig(gotConfig); err != nil { - t.Fatalf("ParseConfig(%+v) failed: %v", gotConfig, err) - } -} - +// TestBuildLeafClusterConfig_DNS tests buildLeafClusterConfig with a +// LOGICAL_DNS cluster configuration, verifying the expected policy hierarchy +// and endpoints. func (s) TestBuildLeafClusterConfig_DNS(t *testing.T) { for _, tt := range []struct { name string @@ -466,6 +470,9 @@ func (s) TestBuildLeafClusterConfig_DNS(t *testing.T) { } } +// TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Disabled tests +// buildLeafClusterConfig with an EDS cluster configuration when pick_first +// weighted shuffling is disabled. func (s) TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, false) @@ -565,6 +572,9 @@ func (s) TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Disabled(t *t } } +// TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Enabled tests +// buildLeafClusterConfig with an EDS cluster configuration when pick_first +// weighted shuffling is enabled. func (s) TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, true) @@ -681,6 +691,9 @@ func (s) TestBuildLeafClusterConfig_EDS_PickFirstWeightedShuffling_Enabled(t *te } } +// TestBuildClusterImplConfigForDNS tests buildClusterImplConfigForDNS +// to verify the generated clusterimpl config and endpoint attributes +// for LOGICAL_DNS clusters. func (s) TestBuildClusterImplConfigForDNS(t *testing.T) { for _, tt := range []struct { name string @@ -754,6 +767,8 @@ func (s) TestBuildClusterImplConfigForDNS(t *testing.T) { } } +// TestBuildClusterImplConfigForEDS_PickFirstWeightedShuffling_Disabled tests +// buildClusterImplConfigForEDS when pick_first weighted shuffling is disabled. func (s) TestBuildClusterImplConfigForEDS_PickFirstWeightedShuffling_Disabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, false) @@ -829,6 +844,8 @@ func (s) TestBuildClusterImplConfigForEDS_PickFirstWeightedShuffling_Disabled(t } } +// TestBuildClusterImplConfigForEDS_PickFirstWeightedShuffling_Enabled tests +// buildClusterImplConfigForEDS when pick_first weighted shuffling is enabled. func (s) TestBuildClusterImplConfigForEDS_PickFirstWeightedShuffling_Enabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, true) @@ -921,6 +938,8 @@ func (s) TestBuildClusterImplConfigForEDS_PickFirstWeightedShuffling_Enabled(t * } } +// TestGroupLocalitiesByPriority tests groupLocalitiesByPriority to ensure +// localities are properly grouped and ordered by priority. func (s) TestGroupLocalitiesByPriority(t *testing.T) { // Create localities for two priorities (p0 and p1). p0Loc0 := makeLocality(0, 20, 0, 2) @@ -989,6 +1008,9 @@ func (s) TestGroupLocalitiesByPriority(t *testing.T) { } } +// TestPriorityLocalitiesToClusterImpl_PickFirstWeightedShuffling_Disabled +// tests priorityLocalitiesToClusterImpl when pick_first weighted shuffling is +// disabled. func (s) TestPriorityLocalitiesToClusterImpl_PickFirstWeightedShuffling_Disabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, false) tests := []struct { @@ -1127,6 +1149,9 @@ func (s) TestPriorityLocalitiesToClusterImpl_PickFirstWeightedShuffling_Disabled } } +// TestPriorityLocalitiesToClusterImpl_PickFirstWeightedShuffling_Enabled tests +// priorityLocalitiesToClusterImpl when pick_first weighted shuffling is +// enabled. func (s) TestPriorityLocalitiesToClusterImpl_PickFirstWeightedShuffling_Enabled(t *testing.T) { testutils.SetEnvConfig(t, &envconfig.PickFirstWeightedShuffling, true) tests := []struct { @@ -1318,6 +1343,9 @@ func testEndpointWithAttrs(endpoint resolver.Endpoint, localityWeight, endpointW return endpoint } +// TestConvertClusterImplMapToOutlierDetection tests +// convertClusterImplMapToOutlierDetection to ensure clusterimpl configs are +// properly wrapped in outlier detection configs. func (s) TestConvertClusterImplMapToOutlierDetection(t *testing.T) { tests := []struct { name string From f2fc6b5cd54b928e44b8df4c19049adce08cb1fc Mon Sep 17 00:00:00 2001 From: Marcin Swierczek Date: Thu, 13 Aug 2026 08:13:02 +0000 Subject: [PATCH 10/10] Revert "Use gracefulswitch.Balancer" This reverts commit 51e986ece5f8d878d1feed588a766573935319a1. --- .../xds/balancer/cdsbalancer/cdsbalancer.go | 62 ++++++++++++++----- 1 file changed, 45 insertions(+), 17 deletions(-) diff --git a/internal/xds/balancer/cdsbalancer/cdsbalancer.go b/internal/xds/balancer/cdsbalancer/cdsbalancer.go index 1e155358898d..9d6496eeee80 100644 --- a/internal/xds/balancer/cdsbalancer/cdsbalancer.go +++ b/internal/xds/balancer/cdsbalancer/cdsbalancer.go @@ -25,7 +25,6 @@ import ( "google.golang.org/grpc/balancer" "google.golang.org/grpc/balancer/base" "google.golang.org/grpc/connectivity" - "google.golang.org/grpc/internal/balancer/gracefulswitch" "google.golang.org/grpc/internal/grpclog" "google.golang.org/grpc/internal/pretty" internalserviceconfig "google.golang.org/grpc/internal/serviceconfig" @@ -40,6 +39,24 @@ import ( const cdsName = "cds_experimental" +var ( + // newChildBalancer is a helper function to build a new child balancer + // and its config parser, and will be overridden in unittests. + newChildBalancer = func(name string, cc balancer.ClientConn, opts balancer.BuildOptions) (balancer.Balancer, balancer.ConfigParser, error) { + builder := balancer.Get(name) + if builder == nil { + return nil, nil, fmt.Errorf("xds: no balancer builder with name %v", name) + } + parser, ok := builder.(balancer.ConfigParser) + if !ok { + return nil, nil, fmt.Errorf("xds: balancer builder for %v does not implement ConfigParser", name) + } + // We directly pass the parent clientConn to the underlying child + // balancer because the cdsBalancer does not deal with subConns. + return builder.Build(cc, opts), parser, nil + } +) + func init() { balancer.Register(bb{}) } @@ -53,7 +70,6 @@ type bb struct{} func (bb) Build(cc balancer.ClientConn, opts balancer.BuildOptions) balancer.Balancer { b := &cdsBalancer{ bOpts: opts, - childLB: gracefulswitch.NewBalancer(cc, opts), clusterConfigs: make(map[string]*xdsresource.ClusterResult), priorityConfigs: make(map[string]*priorityConfig), cc: cc, @@ -95,16 +111,18 @@ type cdsBalancer struct { // The following fields are initialized at build time and are either // read-only after that or provide their own synchronization, and therefore // do not need to be guarded by a mutex. - cc balancer.ClientConn // ClientConn interface passed to child LB. - bOpts balancer.BuildOptions // BuildOptions passed to child LB. - logger *grpclog.PrefixLogger // Prefix logger for all logging. + cc balancer.ClientConn // ClientConn interface passed to child LB. + bOpts balancer.BuildOptions // BuildOptions passed to child LB. + childConfigParser balancer.ConfigParser // Config parser for cluster_resolver LB policy. + logger *grpclog.PrefixLogger // Prefix logger for all logging. // All fields below are accessed only from methods implementing the // balancer.Balancer interface. Since gRPC guarantees that these methods are // never invoked concurrently, no additional synchronization is required to // protect access to these fields. xdsClient xdsclient.XDSClient - childLB *gracefulswitch.Balancer // Graceful switch child policy. + childLB balancer.Balancer // Child policy, built upon resolution of the cluster graph. + childLBName string // Name of the child policy. clusterConfigs map[string]*xdsresource.ClusterResult // Cluster name to the last received result for that cluster. priorityConfigs map[string]*priorityConfig // Hostname to priority config for that leaf cluster. lbCfg *lbConfig // Current load balancing configuration. @@ -242,34 +260,44 @@ func (b *cdsBalancer) updateChildConfig() error { clusterConfig := b.clusterConfigs[clusterName].Config isAggregate := clusterConfig.Cluster.ClusterType == xdsresource.ClusterTypeAggregate + var topLBName string + if isAggregate { + topLBName = priority.Name + } else { + topLBName = outlierdetection.Name + } + + if b.childLB != nil && b.childLBName != topLBName { + b.childLB.Close() + b.childLB = nil + } + if b.childLB == nil { - b.childLB = gracefulswitch.NewBalancer(b.cc, b.bOpts) + childLB, parser, err := newChildBalancer(topLBName, b.cc, b.bOpts) + if err != nil { + return fmt.Errorf("failed to create child policy of type %s: %v", topLBName, err) + } + b.childLB = childLB + b.childLBName = topLBName + b.childConfigParser = parser } var childCfgBytes []byte var endpoints []resolver.Endpoint var err error - var topLBName string if isAggregate { - topLBName = priority.Name childCfgBytes, endpoints, err = buildAggregateClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } else { - topLBName = outlierdetection.Name childCfgBytes, endpoints, err = buildLeafClusterConfigJSON(b.priorities, &b.xdsLBPolicy) } if err != nil { return fmt.Errorf("failed to build child policy config: %v", err) } - cfgJSON, err := json.Marshal([]map[string]json.RawMessage{{topLBName: childCfgBytes}}) - if err != nil { - return fmt.Errorf("failed to marshal child policy config wrapper: %v", err) - } - - childCfg, err := gracefulswitch.ParseConfig(cfgJSON) + childCfg, err := b.childConfigParser.ParseConfig(childCfgBytes) if err != nil { - return fmt.Errorf("failed to parse child policy config: %v", err) + return fmt.Errorf("failed to parse child policy config. This should never happen because the config was generated: %v", err) } if b.logger.V(2) { b.logger.Infof("Built child policy config: %s", pretty.ToJSON(childCfg))