Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
159 changes: 101 additions & 58 deletions internal/provider/kubernetes/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -85,10 +85,13 @@ type gatewayAPIReconciler struct {
eepCRDExists bool
epCRDExists bool
eppCRDExists bool
grpcRouteCRDExists bool
Comment thread
zhaohuabing marked this conversation as resolved.
hrfCRDExists bool
listenerSetCRDExists bool
serviceImportCRDExists bool
spCRDExists bool
tcpRouteCRDExists bool
tlsRouteCRDExists bool
udpRouteCRDExists bool

clusterTrustBundleExits bool
Expand Down Expand Up @@ -1887,28 +1890,34 @@ func (r *gatewayAPIReconciler) processGateways(ctx context.Context, managedGC *g
gtwNamespacedName := utils.NamespacedName(gtw).String()

// ListenerSet Processing (must be done before route processing)
if err := r.processListenerSets(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
if isTransientError(err) {
return err
if r.listenerSetCRDExists {
if err := r.processListenerSets(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
if isTransientError(err) {
return err
}
r.log.Error(err, "failed to process ListenerSets for gateway", "namespace", gtw.Namespace, "name", gtw.Name)
}
r.log.Error(err, "failed to process ListenerSets for gateway", "namespace", gtw.Namespace, "name", gtw.Name)
}

// Route Processing

// Get TLSRoute objects and check if it exists.
if err := r.processTLSRoutes(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
return err
if r.tlsRouteCRDExists {
// Get TLSRoute objects and check if it exists.
if err := r.processTLSRoutes(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
return err
}
}

// Get HTTPRoute objects and check if it exists.
if err := r.processHTTPRoutes(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
return err
}

// Get GRPCRoute objects and check if it exists.
if err := r.processGRPCRoutes(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
return err
if r.grpcRouteCRDExists {
// Get GRPCRoute objects and check if it exists.
if err := r.processGRPCRoutes(ctx, gtwNamespacedName, resourceMap, resourceTree); err != nil {
return err
}
}

if r.tcpRouteCRDExists {
Expand Down Expand Up @@ -2310,25 +2319,44 @@ func (r *gatewayAPIReconciler) watchResources(ctx context.Context, mgr manager.M
return err
}

xlsPredicates := []predicate.TypedPredicate[*gwapiv1.ListenerSet]{
predicate.TypedGenerationChangedPredicate[*gwapiv1.ListenerSet]{},
}
if r.namespaceLabel != nil {
xlsPredicates = append(xlsPredicates, predicate.NewTypedPredicateFuncs(func(obj *gwapiv1.ListenerSet) bool {
return r.hasMatchingNamespaceLabels(obj)
}))
}

if err := c.Watch(
source.Kind(mgr.GetCache(), &gwapiv1.ListenerSet{},
handler.TypedEnqueueRequestsFromMapFunc(func(ctx context.Context, obj *gwapiv1.ListenerSet) []reconcile.Request {
return r.enqueueClass(ctx, obj)
}),
xlsPredicates...)); err != nil {
// The watches for ListenerSet, GRPCRoute, TLSRoute, TCPRoute and UDPRoute are all optional. These
// kinds are in the Gateway API standard channel today (ListenerSet and TLSRoute graduated in
// v1.5, TCPRoute and UDPRoute in v1.6), but they are still not guaranteed to be installed:
// - The cluster may run an older Gateway API bundle, or a standard channel bundle from before
// the kind graduated, so the kind is simply not there.
// - Some managed Kubernetes offerings install only a curated subset of the standard channel and
// don't let users add the missing CRDs. For example, GKE's managed gateway-api-crds addon
// ships Gateway API v1.5 without ListenerSet and GRPCRoute, because the GKE Gateway
// controller doesn't implement them.
// Envoy Gateway must start on those clusters instead of crash-looping: an absent kind can't be
// used there anyway, so its absence should disable the feature, not take down the controller.
r.listenerSetCRDExists, err = checkCRD(resource.KindListenerSet, gwapiv1.GroupVersion.String())
if err != nil {
return err
}
if err := addListenerSetIndexers(ctx, mgr); err != nil {
return err
if !r.listenerSetCRDExists {
r.log.Info("ListenerSet CRD not found, skipping ListenerSet watch")
} else {
xlsPredicates := []predicate.TypedPredicate[*gwapiv1.ListenerSet]{
predicate.TypedGenerationChangedPredicate[*gwapiv1.ListenerSet]{},
}
if r.namespaceLabel != nil {
xlsPredicates = append(xlsPredicates, predicate.NewTypedPredicateFuncs(func(obj *gwapiv1.ListenerSet) bool {
return r.hasMatchingNamespaceLabels(obj)
}))
}

if err := c.Watch(
source.Kind(mgr.GetCache(), &gwapiv1.ListenerSet{},
handler.TypedEnqueueRequestsFromMapFunc(func(ctx context.Context, obj *gwapiv1.ListenerSet) []reconcile.Request {
return r.enqueueClass(ctx, obj)
}),
xlsPredicates...)); err != nil {
return err
}
if err := addListenerSetIndexers(ctx, mgr); err != nil {
return err
}
}

// Watch HTTPRoute CRUDs and process affected Gateways.
Expand All @@ -2350,42 +2378,57 @@ func (r *gatewayAPIReconciler) watchResources(ctx context.Context, mgr manager.M
return err
}

// Watch GRPCRoute CRUDs and process affected Gateways.
grpcrPredicates := commonPredicates[*gwapiv1.GRPCRoute]()
if r.namespaceLabel != nil {
grpcrPredicates = append(grpcrPredicates, predicate.NewTypedPredicateFuncs(func(grpc *gwapiv1.GRPCRoute) bool {
return r.hasMatchingNamespaceLabels(grpc)
}))
}
if err := c.Watch(
source.Kind(mgr.GetCache(), &gwapiv1.GRPCRoute{},
handler.TypedEnqueueRequestsFromMapFunc(func(ctx context.Context, route *gwapiv1.GRPCRoute) []reconcile.Request {
return r.enqueueClass(ctx, route)
}),
grpcrPredicates...)); err != nil {
return err
}
if err := addGRPCRouteIndexers(ctx, mgr); err != nil {
r.grpcRouteCRDExists, err = checkCRD(resource.KindGRPCRoute, gwapiv1.GroupVersion.String())
if err != nil {
return err
}

// Watch TLSRoute CRUDs and process affected Gateways.
tlsrPredicates := commonPredicates[*gwapiv1.TLSRoute]()
if r.namespaceLabel != nil {
tlsrPredicates = append(tlsrPredicates, predicate.NewTypedPredicateFuncs(func(route *gwapiv1.TLSRoute) bool {
return r.hasMatchingNamespaceLabels(route)
}))
if !r.grpcRouteCRDExists {
r.log.Info("GRPCRoute CRD not found, skipping GRPCRoute watch")
} else {
// Watch GRPCRoute CRUDs and process affected Gateways.
grpcrPredicates := commonPredicates[*gwapiv1.GRPCRoute]()
if r.namespaceLabel != nil {
grpcrPredicates = append(grpcrPredicates, predicate.NewTypedPredicateFuncs(func(grpc *gwapiv1.GRPCRoute) bool {
return r.hasMatchingNamespaceLabels(grpc)
}))
}
if err := c.Watch(
source.Kind(mgr.GetCache(), &gwapiv1.GRPCRoute{},
handler.TypedEnqueueRequestsFromMapFunc(func(ctx context.Context, route *gwapiv1.GRPCRoute) []reconcile.Request {
return r.enqueueClass(ctx, route)
}),
grpcrPredicates...)); err != nil {
return err
}
if err := addGRPCRouteIndexers(ctx, mgr); err != nil {
return err
}
}
if err := c.Watch(
source.Kind(mgr.GetCache(), &gwapiv1.TLSRoute{},
handler.TypedEnqueueRequestsFromMapFunc(func(ctx context.Context, route *gwapiv1.TLSRoute) []reconcile.Request {
return r.enqueueClass(ctx, route)
}),
tlsrPredicates...)); err != nil {
r.tlsRouteCRDExists, err = checkCRD(resource.KindTLSRoute, gwapiv1.GroupVersion.String())
if err != nil {
return err
}
if err := addTLSRouteIndexers(ctx, mgr); err != nil {
return err
if !r.tlsRouteCRDExists {
r.log.Info("TLSRoute CRD not found, skipping TLSRoute watch")
} else {
// Watch TLSRoute CRUDs and process affected Gateways.
tlsrPredicates := commonPredicates[*gwapiv1.TLSRoute]()
if r.namespaceLabel != nil {
tlsrPredicates = append(tlsrPredicates, predicate.NewTypedPredicateFuncs(func(route *gwapiv1.TLSRoute) bool {
return r.hasMatchingNamespaceLabels(route)
}))
}
if err := c.Watch(
source.Kind(mgr.GetCache(), &gwapiv1.TLSRoute{},
handler.TypedEnqueueRequestsFromMapFunc(func(ctx context.Context, route *gwapiv1.TLSRoute) []reconcile.Request {
return r.enqueueClass(ctx, route)
}),
tlsrPredicates...)); err != nil {
return err
}
if err := addTLSRouteIndexers(ctx, mgr); err != nil {
return err
}
}

r.udpRouteCRDExists, err = checkCRD(resource.KindUDPRoute, gwapiv1.GroupVersion.String())
Expand Down
3 changes: 3 additions & 0 deletions internal/provider/kubernetes/controller_offline.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,10 +107,13 @@ func NewOfflineGatewayAPIController(
eepCRDExists: true,
epCRDExists: true,
eppCRDExists: true,
grpcRouteCRDExists: true,
hrfCRDExists: true,
listenerSetCRDExists: true,
serviceImportCRDExists: true,
spCRDExists: true,
tcpRouteCRDExists: true,
tlsRouteCRDExists: true,
udpRouteCRDExists: true,
backendCRDExists: true,
}
Expand Down
44 changes: 26 additions & 18 deletions internal/provider/kubernetes/predicates.go
Original file line number Diff line number Diff line change
Expand Up @@ -562,26 +562,30 @@ func (r *gatewayAPIReconciler) isRouteReferencingBackend(nsName *types.Namespace
return true
}

grpcRouteList := &gwapiv1.GRPCRouteList{}
if err := r.client.List(ctx, grpcRouteList, &client.ListOptions{
FieldSelector: fields.OneTermEqualSelector(backendGRPCRouteIndex, nsName.String()),
}); err != nil && !kerrors.IsNotFound(err) {
r.log.Error(err, "failed to find associated GRPCRoutes")
return false
}
if len(grpcRouteList.Items) > 0 {
return true
if r.grpcRouteCRDExists {
grpcRouteList := &gwapiv1.GRPCRouteList{}
if err := r.client.List(ctx, grpcRouteList, &client.ListOptions{
FieldSelector: fields.OneTermEqualSelector(backendGRPCRouteIndex, nsName.String()),
}); err != nil && !kerrors.IsNotFound(err) {
r.log.Error(err, "failed to find associated GRPCRoutes")
return false
}
if len(grpcRouteList.Items) > 0 {
return true
}
}

tlsRouteList := &gwapiv1.TLSRouteList{}
if err := r.client.List(ctx, tlsRouteList, &client.ListOptions{
FieldSelector: fields.OneTermEqualSelector(backendTLSRouteIndex, nsName.String()),
}); err != nil && !kerrors.IsNotFound(err) {
r.log.Error(err, "failed to find associated TLSRoutes")
return false
}
if len(tlsRouteList.Items) > 0 {
return true
if r.tlsRouteCRDExists {
tlsRouteList := &gwapiv1.TLSRouteList{}
if err := r.client.List(ctx, tlsRouteList, &client.ListOptions{
FieldSelector: fields.OneTermEqualSelector(backendTLSRouteIndex, nsName.String()),
}); err != nil && !kerrors.IsNotFound(err) {
r.log.Error(err, "failed to find associated TLSRoutes")
return false
}
if len(tlsRouteList.Items) > 0 {
return true
}
}

if r.tcpRouteCRDExists {
Expand Down Expand Up @@ -1007,6 +1011,10 @@ func (r *gatewayAPIReconciler) isRouteReferencingHTTPRouteFilter(nsName *types.N
return true
}

if !r.grpcRouteCRDExists {
return false
}

grpcRouteList := &gwapiv1.GRPCRouteList{}
if err := r.client.List(ctx, grpcRouteList, &client.ListOptions{
FieldSelector: fields.OneTermEqualSelector(httpRouteFilterGRPCRouteIndex, nsName.String()),
Expand Down
63 changes: 55 additions & 8 deletions internal/provider/kubernetes/predicates_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1154,10 +1154,15 @@ func TestValidateServiceForReconcile(t *testing.T) {
sampleServiceBackendRef := test.GetServiceBackendRef(types.NamespacedName{Name: "service"}, 80)

testCases := []struct {
name string
configs []client.Object
service client.Object
expect bool
name string
// grpcRouteCRDAbsent and tlsRouteCRDAbsent simulate a cluster whose Gateway API
// CRD bundle doesn't include GRPCRoute or TLSRoute, e.g. GKE's managed
// gateway-api-crds addon.
grpcRouteCRDAbsent bool
tlsRouteCRDAbsent bool
configs []client.Object
service client.Object
expect bool
}{
{
name: "gateway service but deployment or daemonset does not exist",
Expand Down Expand Up @@ -1242,6 +1247,17 @@ func TestValidateServiceForReconcile(t *testing.T) {
service: test.GetService(types.NamespacedName{Name: "service"}, nil, nil),
expect: true,
},
{
name: "grpc route service routes exist but GRPCRoute CRD is absent",
grpcRouteCRDAbsent: true,
configs: []client.Object{
test.GetGatewayClass("test-gc", egv1a1.GatewayControllerName, nil),
sampleGateway,
test.GetGRPCRoute(types.NamespacedName{Name: "grpcroute-test"}, "scheduled-status-test", types.NamespacedName{Name: "service"}, 80),
},
service: test.GetService(types.NamespacedName{Name: "service"}, nil, nil),
expect: false,
},
{
name: "tls route service routes exist",
configs: []client.Object{
Expand All @@ -1253,6 +1269,18 @@ func TestValidateServiceForReconcile(t *testing.T) {
service: test.GetService(types.NamespacedName{Name: "service"}, nil, nil),
expect: true,
},
{
name: "tls route service routes exist but TLSRoute CRD is absent",
tlsRouteCRDAbsent: true,
configs: []client.Object{
test.GetGatewayClass("test-gc", egv1a1.GatewayControllerName, nil),
sampleGateway,
test.GetTLSRoute(types.NamespacedName{Name: "tlsroute-test"}, "scheduled-status-test",
types.NamespacedName{Name: "service"}, 443),
},
service: test.GetService(types.NamespacedName{Name: "service"}, nil, nil),
expect: false,
},
{
name: "udp route service routes exist",
configs: []client.Object{
Expand Down Expand Up @@ -1496,6 +1524,8 @@ func TestValidateServiceForReconcile(t *testing.T) {
}

for _, tc := range testCases {
r.grpcRouteCRDExists = !tc.grpcRouteCRDAbsent
r.tlsRouteCRDExists = !tc.tlsRouteCRDAbsent
r.client = fakeclient.NewClientBuilder().
WithScheme(envoygateway.GetScheme()).
WithObjects(tc.configs...).
Expand Down Expand Up @@ -1846,10 +1876,13 @@ func TestValidateHTTPRouteFilerForReconcile(t *testing.T) {
sampleHTTPRouteFilter := test.GetHTTPRouteFilter(types.NamespacedName{Name: "httproutefilter"})

testCases := []struct {
name string
configs []client.Object
httpRouteFilter client.Object
expect bool
name string
// grpcRouteCRDAbsent simulates a cluster whose Gateway API CRD bundle doesn't
// include GRPCRoute, e.g. GKE's managed gateway-api-crds addon.
grpcRouteCRDAbsent bool
configs []client.Object
httpRouteFilter client.Object
expect bool
}{
{
name: "httproutefilter but not referenced by route",
Expand Down Expand Up @@ -1886,6 +1919,19 @@ func TestValidateHTTPRouteFilerForReconcile(t *testing.T) {
httpRouteFilter: sampleHTTPRouteFilter,
expect: true,
},
{
name: "httproutefilter referenced by grpcroute but GRPCRoute CRD is absent",
grpcRouteCRDAbsent: true,
configs: []client.Object{
sampleGWC,
sampleGateway,
sampleService,
sampleHTTPRouteFilter,
test.GetGRPCRouteWithHTTPRouteFilter(types.NamespacedName{Name: "grpcroute-test"}, "scheduled-status-test", types.NamespacedName{Name: "service"}, 80, "httproutefilter"),
},
httpRouteFilter: sampleHTTPRouteFilter,
expect: false,
},
}

// Create the reconciler.
Expand All @@ -1897,6 +1943,7 @@ func TestValidateHTTPRouteFilerForReconcile(t *testing.T) {
}

for _, tc := range testCases {
r.grpcRouteCRDExists = !tc.grpcRouteCRDAbsent
r.client = fakeclient.NewClientBuilder().
WithScheme(envoygateway.GetScheme()).
WithObjects(tc.configs...).
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Fixed Envoy Gateway crash-looping at startup on clusters whose Gateway API CRD bundle omits ListenerSet, GRPCRoute or TLSRoute, such as GKE's managed gateway-api-crds addon, by making those watches conditional on the CRD being present.