Skip to content
386 changes: 386 additions & 0 deletions internal/xds/balancer/cdsbalancer/e2e_test/aggregate_cluster_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,11 @@ import (
v3listenerpb "github.com/envoyproxy/go-control-plane/envoy/config/listener/v3"
v3routepb "github.com/envoyproxy/go-control-plane/envoy/config/route/v3"
v3discoverypb "github.com/envoyproxy/go-control-plane/envoy/service/discovery/v3"
v3lrspb "github.com/envoyproxy/go-control-plane/envoy/service/load_stats/v3"
Comment thread
Pranjali-2501 marked this conversation as resolved.
"google.golang.org/grpc/internal/testutils/xds/fakeserver"
testgrpc "google.golang.org/grpc/interop/grpc_testing"
testpb "google.golang.org/grpc/interop/grpc_testing"
"google.golang.org/protobuf/types/known/durationpb"

_ "google.golang.org/grpc/internal/xds/httpfilter/router" // Register the router filter
)
Expand Down Expand Up @@ -1282,3 +1285,386 @@ func addrsToEndpoints(addrs []resolver.Address) []resolver.Endpoint {
}
return endpoints
}

// TestAggregateCluster_LRS_PrimaryLeafOnly tests that when an aggregate cluster
// (Root -> [PrimaryEDS, SecondaryEDS]) has LRS enabled on both leaf clusters,
// LRS reports are sent only for the primary leaf cluster while it is healthy.
func (s) TestAggregateCluster_LRS_PrimaryLeafOnly(t *testing.T) {
Comment thread
Pranjali-2501 marked this conversation as resolved.
Outdated
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
defer cancel()

const clusterName1 = clusterName + "-cluster-1"
const clusterName2 = clusterName + "-cluster-2"

managementServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{
AllowResourceSubset: true,
SupportLoadReportingService: true,
})
nodeID := uuid.New().String()
bootstrapContents := e2e.DefaultBootstrapContents(t, nodeID, managementServer.Address)

servers, cleanup2 := startTestServiceBackends(t, 2)
defer cleanup2()
_, ports := backendAddressesAndPorts(t, servers)

resources := e2e.UpdateOptions{
NodeID: nodeID,
Listeners: []*v3listenerpb.Listener{e2e.DefaultClientListener(serviceName, routeName)},
Routes: []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(routeName, serviceName, clusterName)},
Clusters: []*v3clusterpb.Cluster{
makeAggregateClusterResource(clusterName, []string{clusterName1, clusterName2}),
e2e.ClusterResourceWithOptions(e2e.ClusterOptions{
ClusterName: clusterName1,
Type: e2e.ClusterTypeEDS,
EnableLRS: true,
}),
e2e.ClusterResourceWithOptions(e2e.ClusterOptions{
ClusterName: clusterName2,
Type: e2e.ClusterTypeEDS,
EnableLRS: true,
}),
},
Endpoints: []*v3endpointpb.ClusterLoadAssignment{
e2e.DefaultEndpoint(clusterName1, "localhost", []uint32{uint32(ports[0])}),
e2e.DefaultEndpoint(clusterName2, "localhost", []uint32{uint32(ports[1])}),
},
SkipValidation: true,
Comment thread
Pranjali-2501 marked this conversation as resolved.
Outdated
}
if err := managementServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}

cc, cleanup := setupAndDial(t, bootstrapContents)
defer cleanup()

client := testgrpc.NewTestServiceClient(cc)

// Ensure LRS stream is opened.
if _, err := managementServer.LRSServer.LRSStreamOpenChan.Receive(ctx); err != nil {
t.Fatalf("Timeout waiting for LRS stream open: %v", err)
}
if _, err := managementServer.LRSServer.LRSRequestChan.Receive(ctx); err != nil {
t.Fatalf("Timeout waiting for initial LRS request: %v", err)
}

// Send LRS response requesting load reporting every 10ms.
managementServer.LRSServer.LRSResponseChan <- &fakeserver.Response{
Resp: &v3lrspb.LoadStatsResponse{
SendAllClusters: true,
LoadReportingInterval: durationpb.New(defaultLoadReportingInterval),
},
}

const wantRPCCount = 10
// Make RPC calls and verify every call reaches the primary cluster (clusterName1).
for i := 0; i < wantRPCCount; i++ {
peer := &peer.Peer{}
if _, err := client.EmptyCall(ctx, &testpb.Empty{}, grpc.Peer(peer), grpc.WaitForReady(true)); err != nil {
t.Fatalf("EmptyCall() failed: %v", err)
}
if got, want := peer.Addr.String(), servers[0].Address; got != want {
t.Fatalf("EmptyCall() call #%d routed to %q, want %q", i, got, want)
}
}

// Verify LRS reports load for clusterName1 and 0 load for clusterName2.
for {
select {
case <-ctx.Done():
t.Fatalf("Timeout waiting for LRS load report for %q", clusterName1)
case req := <-managementServer.LRSServer.LRSRequestChan.C:
loadStats := req.(*fakeserver.Request).Req.(*v3lrspb.LoadStatsRequest)
for _, load := range loadStats.ClusterStats {
// For unused leaf clusters (when SendAllClusters: true), LRS reports are
// expected to be empty (0 successful requests across any locality).
if load.ClusterName == clusterName2 {
for _, loc := range load.UpstreamLocalityStats {
if loc.TotalSuccessfulRequests > 0 {
t.Fatalf("Received unexpected LRS load report for secondary cluster %q: %+v", clusterName2, load)
}
}
}
if load.ClusterName == clusterName1 {
// Each leaf cluster has exactly one locality.
if len(load.UpstreamLocalityStats) != 1 {
t.Fatalf("UpstreamLocalityStats length = %d, want 1", len(load.UpstreamLocalityStats))
}
if load.UpstreamLocalityStats[0].TotalSuccessfulRequests == wantRPCCount {
return
}
}
if load.ClusterName != clusterName1 && load.ClusterName != clusterName2 {
t.Fatalf("Unexpected cluster name %q", load.ClusterName)
}
}
}
}
}

// TestAggregateCluster_LRS_FailoverToSecondaryLeaf tests that when the primary
// leaf cluster fails, traffic shifts to the secondary leaf cluster and LRS stats
// are reported for the secondary leaf cluster.
func (s) TestAggregateCluster_LRS_FailoverToSecondaryLeaf(t *testing.T) {
Comment thread
Pranjali-2501 marked this conversation as resolved.
Outdated
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
defer cancel()

const clusterName1 = clusterName + "-cluster-1"
const clusterName2 = clusterName + "-cluster-2"

managementServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{
AllowResourceSubset: true,
SupportLoadReportingService: true,
})
nodeID := uuid.New().String()
bootstrapContents := e2e.DefaultBootstrapContents(t, nodeID, managementServer.Address)

servers, cleanup2 := startTestServiceBackends(t, 2)
defer cleanup2()
_, ports := backendAddressesAndPorts(t, servers)

resources := e2e.UpdateOptions{
NodeID: nodeID,
Listeners: []*v3listenerpb.Listener{e2e.DefaultClientListener(serviceName, routeName)},
Routes: []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(routeName, serviceName, clusterName)},
Clusters: []*v3clusterpb.Cluster{
makeAggregateClusterResource(clusterName, []string{clusterName1, clusterName2}),
e2e.ClusterResourceWithOptions(e2e.ClusterOptions{
ClusterName: clusterName1,
Type: e2e.ClusterTypeEDS,
EnableLRS: true,
}),
e2e.ClusterResourceWithOptions(e2e.ClusterOptions{
ClusterName: clusterName2,
Type: e2e.ClusterTypeEDS,
EnableLRS: true,
}),
},
Endpoints: []*v3endpointpb.ClusterLoadAssignment{
e2e.DefaultEndpoint(clusterName1, "localhost", []uint32{uint32(ports[0])}),
e2e.DefaultEndpoint(clusterName2, "localhost", []uint32{uint32(ports[1])}),
},
SkipValidation: true,
}
if err := managementServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}

cc, cleanup := setupAndDial(t, bootstrapContents)
defer cleanup()

client := testgrpc.NewTestServiceClient(cc)

// Ensure LRS stream is opened.
if _, err := managementServer.LRSServer.LRSStreamOpenChan.Receive(ctx); err != nil {
t.Fatalf("Timeout waiting for LRS stream open: %v", err)
}
if _, err := managementServer.LRSServer.LRSRequestChan.Receive(ctx); err != nil {
t.Fatalf("Timeout waiting for initial LRS request: %v", err)
}

managementServer.LRSServer.LRSResponseChan <- &fakeserver.Response{
Resp: &v3lrspb.LoadStatsResponse{
SendAllClusters: true,
LoadReportingInterval: durationpb.New(defaultLoadReportingInterval),
},
}

// Make RPCs to primary backend.
peer := &peer.Peer{}
if _, err := client.EmptyCall(ctx, &testpb.Empty{}, grpc.Peer(peer), grpc.WaitForReady(true)); err != nil {
t.Fatalf("EmptyCall() failed: %v", err)
}
if got, want := peer.Addr.String(), servers[0].Address; got != want {
t.Fatalf("EmptyCall() routed to %q, want %q", got, want)
}

const wantRPCCount = 1

// Verify LRS report includes load for primary cluster (clusterName1).
for gotCluster1Report := false; !gotCluster1Report; {
select {
case <-ctx.Done():
t.Fatalf("Timeout waiting for LRS load report for primary cluster %q", clusterName1)
case req := <-managementServer.LRSServer.LRSRequestChan.C:
loadStats := req.(*fakeserver.Request).Req.(*v3lrspb.LoadStatsRequest)
for _, load := range loadStats.ClusterStats {
if load.ClusterName == clusterName1 {
if len(load.UpstreamLocalityStats) != 1 {
t.Fatalf("UpstreamLocalityStats length = %d, want 1", len(load.UpstreamLocalityStats))
}
if load.UpstreamLocalityStats[0].TotalSuccessfulRequests == wantRPCCount {
gotCluster1Report = true
}
if load.UpstreamLocalityStats[0].TotalSuccessfulRequests > wantRPCCount {
t.Fatalf("Too many RPCs to primary cluster %q", clusterName1)
}
}
if load.ClusterName == clusterName2 {
if len(load.UpstreamLocalityStats) != 1 {
t.Fatalf("UpstreamLocalityStats length = %d, want 1", len(load.UpstreamLocalityStats))
}
if load.UpstreamLocalityStats[0].TotalSuccessfulRequests > 0 {
t.Fatalf("Received unexpected LRS load report for secondary cluster %q: %+v", clusterName2, load)
}
}
if load.ClusterName != clusterName1 && load.ClusterName != clusterName2 {
t.Fatalf("Unexpected cluster name %q", load.ClusterName)
}
}
}
}

// Fail primary cluster by clearing endpoints for clusterName1.
resources.Endpoints = []*v3endpointpb.ClusterLoadAssignment{
e2e.DefaultEndpoint(clusterName1, "localhost", nil),
e2e.DefaultEndpoint(clusterName2, "localhost", []uint32{uint32(ports[1])}),
}
if err := managementServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}

// Make RPCs and verify failover to secondary backend.
for ctx.Err() == nil {
Comment thread
easwars marked this conversation as resolved.
Outdated
if _, err := client.EmptyCall(ctx, &testpb.Empty{}, grpc.Peer(peer)); err == nil && peer.Addr.String() == servers[1].Address {
break
}
}
if ctx.Err() != nil {
t.Fatalf("Timeout waiting for RPCs to switch to secondary backend %q", servers[1].Address)
}

// Verify LRS report now includes load for secondary cluster (clusterName2).
for {
select {
case <-ctx.Done():
t.Fatalf("Timeout waiting for LRS load report for secondary cluster %q", clusterName2)
case req := <-managementServer.LRSServer.LRSRequestChan.C:
loadStats := req.(*fakeserver.Request).Req.(*v3lrspb.LoadStatsRequest)
Comment thread
easwars marked this conversation as resolved.
for _, load := range loadStats.ClusterStats {
if load.ClusterName == clusterName2 {
if len(load.UpstreamLocalityStats) != 1 {
t.Fatalf("UpstreamLocalityStats length = %d, want 1", len(load.UpstreamLocalityStats))
}
if load.UpstreamLocalityStats[0].TotalSuccessfulRequests == wantRPCCount {
return
}
if load.UpstreamLocalityStats[0].TotalSuccessfulRequests > wantRPCCount {
t.Fatalf("Too many RPCs to secondary cluster %q", clusterName2)
}
}
if load.ClusterName != clusterName1 && load.ClusterName != clusterName2 {
t.Fatalf("Unexpected cluster name %q", load.ClusterName)
}
}
}
}
}

// TestAggregateCluster_LRS_DropsAndLoadReporting tests that overload drops
// configured on a leaf cluster in an aggregate cluster hierarchy are correctly
// reported to LRS.
func (s) TestAggregateCluster_LRS_DropsAndLoadReporting(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
defer cancel()

const clusterName1 = clusterName + "-cluster-1"

managementServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{
AllowResourceSubset: true,
SupportLoadReportingService: true,
})
nodeID := uuid.New().String()
bootstrapContents := e2e.DefaultBootstrapContents(t, nodeID, managementServer.Address)

servers, cleanup2 := startTestServiceBackends(t, 1)
defer cleanup2()
_, ports := backendAddressesAndPorts(t, servers)

const dropCategory = "test_drop"
resources := e2e.UpdateOptions{
NodeID: nodeID,
Listeners: []*v3listenerpb.Listener{e2e.DefaultClientListener(serviceName, routeName)},
Routes: []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(routeName, serviceName, clusterName)},
Clusters: []*v3clusterpb.Cluster{
makeAggregateClusterResource(clusterName, []string{clusterName1}),
e2e.ClusterResourceWithOptions(e2e.ClusterOptions{
ClusterName: clusterName1,
Type: e2e.ClusterTypeEDS,
EnableLRS: true,
}),
},
Endpoints: []*v3endpointpb.ClusterLoadAssignment{
e2e.EndpointResourceWithOptions(e2e.EndpointOptions{
ClusterName: clusterName1,
Host: "localhost",
Localities: []e2e.LocalityOptions{
{
Backends: []e2e.BackendOptions{{Ports: []uint32{uint32(ports[0])}, Weight: 1}},
Weight: 1,
},
},
DropPercents: map[string]int{dropCategory: 100},
}),
},
SkipValidation: true,
}
if err := managementServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}

cc, cleanup := setupAndDial(t, bootstrapContents)
defer cleanup()

client := testgrpc.NewTestServiceClient(cc)

if _, err := managementServer.LRSServer.LRSStreamOpenChan.Receive(ctx); err != nil {
t.Fatalf("Timeout waiting for LRS stream open: %v", err)
}
if _, err := managementServer.LRSServer.LRSRequestChan.Receive(ctx); err != nil {
t.Fatalf("Timeout waiting for initial LRS request: %v", err)
}

managementServer.LRSServer.LRSResponseChan <- &fakeserver.Response{
Resp: &v3lrspb.LoadStatsResponse{
SendAllClusters: true,
LoadReportingInterval: durationpb.New(defaultLoadReportingInterval),
},
}

const wantRPCCount = 5
// Make RPC calls; all should be dropped locally by xDS overload drop.
for i := 0; i < wantRPCCount; i++ {
_, err := client.EmptyCall(ctx, &testpb.Empty{}, grpc.WaitForReady(true))
if status.Code(err) != codes.Unavailable || !strings.Contains(err.Error(), "RPC is dropped") {
t.Fatalf("EmptyCall() error = %v, want Unavailable with message 'RPC is dropped'", err)
}
}

// Verify LRS report contains drop statistics for dropCategory.
for {
select {
case <-ctx.Done():
t.Fatalf("Timeout waiting for LRS drop report for %q", clusterName1)
case req := <-managementServer.LRSServer.LRSRequestChan.C:
loadStats := req.(*fakeserver.Request).Req.(*v3lrspb.LoadStatsRequest)
for _, load := range loadStats.ClusterStats {
if load.ClusterName == clusterName1 {
for _, drop := range load.DroppedRequests {
if drop.Category == dropCategory && drop.DroppedCount == wantRPCCount {
return
}
if drop.Category == dropCategory && drop.DroppedCount > wantRPCCount {
t.Fatalf("Too many RPCs dropped for category %q", dropCategory)
}
if drop.Category != dropCategory {
t.Fatalf("Unexpected drop category %q", drop.Category)
}
}
}
if load.ClusterName != clusterName1 {
t.Fatalf("Unexpected cluster name %q", load.ClusterName)
}
}
}
}
}
Loading
Loading