Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 29 additions & 4 deletions internal/xds/resolver/helpers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,8 @@ import (
"google.golang.org/grpc/serviceconfig"
"google.golang.org/grpc/status"

v3clusterpb "github.com/envoyproxy/go-control-plane/envoy/config/cluster/v3"
v3endpointpb "github.com/envoyproxy/go-control-plane/envoy/config/endpoint/v3"
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"
Expand All @@ -60,11 +62,17 @@ const (
defaultTestServiceName = "service-name"
defaultTestRouteConfigName = "route-config-name"
defaultTestClusterName = "cluster-name"
defaultTestEndpointName = "endpoint-name"
defaultTestHostname = "test-host"
)

// This is the expected service config when using default listener and route
// configuration resources from the e2e package using the above resource names.
var wantDefaultServiceConfig = fmt.Sprintf(`{
var defaultTestPort = []uint32{8080}

// wantServiceConfig returns a JSON representation of a service config with
// xds_cluster_manager_experimental LB policy with a child policy of
// cds_experimental for the provided cluster name.
func wantServiceConfig(clusterName string) string {
return fmt.Sprintf(`{
"loadBalancingConfig": [{
"xds_cluster_manager_experimental": {
"children": {
Expand All @@ -78,7 +86,8 @@ var wantDefaultServiceConfig = fmt.Sprintf(`{
}
}
}]
}`, defaultTestClusterName, defaultTestClusterName)
}`, clusterName, clusterName)
}

// buildResolverForTarget builds an xDS resolver for the given target. If
// the bootstrap contents are provided, it build the xDS resolver using them
Expand Down Expand Up @@ -283,6 +292,22 @@ func configureResourcesOnManagementServer(ctx context.Context, t *testing.T, mgm
}
}

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we need to skip validation here? We do that in configureResourcesOnManagementServer because we don't specify all resources there.

}
if err := mgmtServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}
}

// waitForResourceNames waits for the wantNames to be pushed on to namesCh.
// Fails the test by calling t.Fatal if the context expires before that.
func waitForResourceNames(ctx context.Context, t *testing.T, namesCh chan []string, wantNames []string) {
Expand Down
67 changes: 40 additions & 27 deletions internal/xds/resolver/watch_service_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,28 +49,35 @@ func (s) TestServiceWatch_ListenerPointsToNewRouteConfiguration(t *testing.T) {
mgmtServer, lisCh, routeCfgCh, bc := setupManagementServerForTest(t, nodeID)

// Configure resources on the management server.
listeners := []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, defaultTestRouteConfigName)}
routes := []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(defaultTestRouteConfigName, defaultTestServiceName, defaultTestClusterName)}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
resources := e2e.DefaultClientResources(e2e.ResourceParams{
DialTarget: defaultTestServiceName,
NodeID: nodeID,
Host: defaultTestHostname,
Port: defaultTestPort[0],
SecLevel: e2e.SecurityLevelNone,
})
if err := mgmtServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}

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

// Verify initial update from the resolver.
waitForResourceNames(ctx, t, lisCh, []string{defaultTestServiceName})
waitForResourceNames(ctx, t, routeCfgCh, []string{defaultTestRouteConfigName})
verifyUpdateFromResolver(ctx, t, stateCh, wantDefaultServiceConfig)
waitForResourceNames(ctx, t, routeCfgCh, []string{resources.Routes[0].Name})
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(resources.Clusters[0].Name))

// Update the listener resource to point to a new route configuration name.
// Leave the old route configuration resource unchanged.
newTestRouteConfigName := defaultTestRouteConfigName + "-new"
listeners = []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, newTestRouteConfigName)}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
resources.Listeners = []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, newTestRouteConfigName)}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, resources.Listeners, resources.Routes)

// Verify that the new route configuration resource is requested.
waitForResourceNames(ctx, t, routeCfgCh, []string{newTestRouteConfigName})

// Update the old route configuration resource by adding a new route.
routes[0].VirtualHosts[0].Routes = append(routes[0].VirtualHosts[0].Routes, &v3routepb.Route{
resources.Routes[0].VirtualHosts[0].Routes = append(resources.Routes[0].VirtualHosts[0].Routes, &v3routepb.Route{
Match: &v3routepb.RouteMatch{
PathSpecifier: &v3routepb.RouteMatch_Prefix{Prefix: "/foo/bar"},
CaseSensitive: &wrapperspb.BoolValue{Value: false},
Expand All @@ -81,17 +88,17 @@ func (s) TestServiceWatch_ListenerPointsToNewRouteConfiguration(t *testing.T) {
},
},
})
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, resources.Listeners, resources.Routes)

// Wait for no update from the resolver.
verifyNoUpdateFromResolver(ctx, t, stateCh)

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

// Ensure update from the resolver.
verifyUpdateFromResolver(ctx, t, stateCh, wantDefaultServiceConfig)
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(resources.Clusters[0].Name))
}

// Tests the case where the listener resource changes to contain an inline route
Expand All @@ -106,22 +113,28 @@ func (s) TestServiceWatch_ListenerPointsToInlineRouteConfiguration(t *testing.T)
mgmtServer, lisCh, routeCfgCh, bc := setupManagementServerForTest(t, nodeID)

// Configure resources on the management server.
listeners := []*v3listenerpb.Listener{e2e.DefaultClientListener(defaultTestServiceName, defaultTestRouteConfigName)}
routes := []*v3routepb.RouteConfiguration{e2e.DefaultRouteConfig(defaultTestRouteConfigName, defaultTestServiceName, defaultTestClusterName)}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, routes)

resources := e2e.DefaultClientResources(e2e.ResourceParams{
DialTarget: defaultTestServiceName,
NodeID: nodeID,
Host: defaultTestHostname,
Port: defaultTestPort[0],
SecLevel: e2e.SecurityLevelNone,
})
if err := mgmtServer.Update(ctx, resources); err != nil {
t.Fatal(err)
}
stateCh, _, _ := buildResolverForTarget(t, resolver.Target{URL: *testutils.MustParseURL("xds:///" + defaultTestServiceName)}, bc)

// Verify initial update from the resolver.
waitForResourceNames(ctx, t, lisCh, []string{defaultTestServiceName})
waitForResourceNames(ctx, t, routeCfgCh, []string{defaultTestRouteConfigName})
verifyUpdateFromResolver(ctx, t, stateCh, wantDefaultServiceConfig)
waitForResourceNames(ctx, t, routeCfgCh, []string{resources.Routes[0].Name})
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(resources.Clusters[0].Name))

// Update listener to contain an inline route configuration.
hcm := testutils.MarshalAny(t, &v3httppb.HttpConnectionManager{
RouteSpecifier: &v3httppb.HttpConnectionManager_RouteConfig{
RouteConfig: &v3routepb.RouteConfiguration{
Name: defaultTestRouteConfigName,
Name: resources.Routes[0].Name,
VirtualHosts: []*v3routepb.VirtualHost{{
Domains: []string{defaultTestServiceName},
Routes: []*v3routepb.Route{{
Expand All @@ -130,7 +143,7 @@ func (s) TestServiceWatch_ListenerPointsToInlineRouteConfiguration(t *testing.T)
},
Action: &v3routepb.Route_Route{
Route: &v3routepb.RouteAction{
ClusterSpecifier: &v3routepb.RouteAction_Cluster{Cluster: defaultTestClusterName},
ClusterSpecifier: &v3routepb.RouteAction_Cluster{Cluster: resources.Clusters[0].Name},
},
},
}},
Expand All @@ -139,7 +152,7 @@ func (s) TestServiceWatch_ListenerPointsToInlineRouteConfiguration(t *testing.T)
},
HttpFilters: []*v3httppb.HttpFilter{e2e.HTTPFilter("router", &v3routerpb.Router{})},
})
listeners = []*v3listenerpb.Listener{{
resources.Listeners = []*v3listenerpb.Listener{{
Name: defaultTestServiceName,
ApiListener: &v3listenerpb.ApiListener{ApiListener: hcm},
FilterChains: []*v3listenerpb.FilterChain{{
Expand All @@ -150,19 +163,19 @@ func (s) TestServiceWatch_ListenerPointsToInlineRouteConfiguration(t *testing.T)
}},
}},
}}
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, listeners, nil)
configureResourcesOnManagementServer(ctx, t, mgmtServer, nodeID, resources.Listeners, nil)

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

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

// Verify that that route configuration resource is requested.
waitForResourceNames(ctx, t, routeCfgCh, []string{defaultTestRouteConfigName})
waitForResourceNames(ctx, t, routeCfgCh, []string{resources.Routes[0].Name})

// Verify that appropriate SC is pushed on the channel.
verifyUpdateFromResolver(ctx, t, stateCh, wantDefaultServiceConfig)
verifyUpdateFromResolver(ctx, t, stateCh, wantServiceConfig(resources.Clusters[0].Name))
}
Loading