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
53 changes: 24 additions & 29 deletions internal/provider/kubernetes/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -605,72 +605,71 @@ func (r *gatewayAPIReconciler) processBackendRefs(ctx context.Context, gwcResour
// The Gateway API translation layer will surface these errors in the status of the resources referencing them.
for backendRef := range resourceMappings.allAssociatedBackendRefs {
backendRefKind := gatewayapi.KindDerefOr(backendRef.Kind, resource.KindService)
r.log.Info("processing Backend", "kind", backendRefKind, "namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))

backendNs, backendName := string(*backendRef.Namespace), string(backendRef.Name)
logger := r.log.WithValues("kind", backendRefKind, "namespace", backendNs,
"name", backendName)
logger.Info("processing Backend")
nn := types.NamespacedName{Namespace: backendNs, Name: backendName}
var endpointSliceLabelKey string
switch backendRefKind {
case resource.KindService:
service := new(corev1.Service)
err := r.client.Get(ctx, types.NamespacedName{Namespace: string(*backendRef.Namespace), Name: string(backendRef.Name)}, service)
err := r.client.Get(ctx, nn, service)
if err != nil {
if isTransientError(err) {
return err
}
r.log.Error(err, "failed to get Service", "namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
logger.Error(err, "failed to get Service")
} else {
resourceMappings.allAssociatedNamespaces.Insert(service.Namespace)
gwcResource.Services = append(gwcResource.Services, service)
r.log.Info("added Service to resource tree", "namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
svcKey := utils.NamespacedName(service).String()
if !resourceMappings.allAssociatedServices.Has(svcKey) {
resourceMappings.allAssociatedServices.Insert(svcKey)
gwcResource.Services = append(gwcResource.Services, service)
logger.Info("added Service to resource tree")
}
Comment on lines +625 to +630

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

this's the core change.

}
endpointSliceLabelKey = discoveryv1.LabelServiceName

case resource.KindServiceImport:
serviceImport := new(mcsapiv1a1.ServiceImport)
err := r.client.Get(ctx, types.NamespacedName{Namespace: string(*backendRef.Namespace), Name: string(backendRef.Name)}, serviceImport)
err := r.client.Get(ctx, nn, serviceImport)
if err != nil {
if isTransientError(err) {
return err
}
r.log.Error(err, "failed to get ServiceImport", "namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
logger.Error(err, "failed to get ServiceImport")
} else {
resourceMappings.allAssociatedNamespaces.Insert(serviceImport.Namespace)
key := utils.NamespacedName(serviceImport).String()
if !resourceMappings.allAssociatedServiceImports.Has(key) {
resourceMappings.allAssociatedServiceImports.Insert(key)
gwcResource.ServiceImports = append(gwcResource.ServiceImports, serviceImport)
r.log.Info("added ServiceImport to resource tree", "namespace", string(*backendRef.Namespace),
logger.Info("added ServiceImport to resource tree", "namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
}
}
endpointSliceLabelKey = mcsapiv1a1.LabelServiceName

case egv1a1.KindBackend:
if r.backendAPIDisabled() {
r.log.V(6).Info("skipping Backend processing as Backend API is disabled.")
logger.V(6).Info("skipping Backend processing as Backend API is disabled.")
continue
}
backend := new(egv1a1.Backend)
err := r.client.Get(ctx, types.NamespacedName{Namespace: string(*backendRef.Namespace), Name: string(backendRef.Name)}, backend)
err := r.client.Get(ctx, nn, backend)
if err != nil {
if isTransientError(err) {
return err
}
r.log.Error(err, "failed to get Backend",
"namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
logger.Error(err, "failed to get Backend")
} else {
resourceMappings.allAssociatedNamespaces.Insert(backend.Namespace)
key := utils.NamespacedName(backend).String()
if !resourceMappings.allAssociatedBackends.Has(key) {
resourceMappings.allAssociatedBackends.Insert(key)
gwcResource.Backends = append(gwcResource.Backends, backend)
r.log.Info("added Backend to resource tree",
"namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
logger.Info("added Backend to resource tree")
}
}

Expand Down Expand Up @@ -702,8 +701,7 @@ func (r *gatewayAPIReconciler) processBackendRefs(ctx context.Context, gwcResour
if isTransientError(err) {
return err
}
r.log.Error(err,
"failed to process CACertificateRef for Backend",
logger.Error(err, "failed to process CACertificateRef for Backend",
"backend", backend, "caCertificateRef", caCertRef.Name)
}
}
Expand All @@ -727,9 +725,7 @@ func (r *gatewayAPIReconciler) processBackendRefs(ctx context.Context, gwcResour
if extRefFilter, exists := resourceMappings.extensionRefFilters[key]; exists {
resourceMappings.extensionRefFilters[key] = extRefFilter
gwcResource.ExtensionRefFilters = append(gwcResource.ExtensionRefFilters, extRefFilter)
r.log.Info("added custom backend resource to resource tree",
"kind", backendRefKind, "namespace", string(*backendRef.Namespace),
"name", string(backendRef.Name))
logger.Info("added custom backend resource to resource tree")
}
}
}
Expand All @@ -748,15 +744,14 @@ func (r *gatewayAPIReconciler) processBackendRefs(ctx context.Context, gwcResour
if isTransientError(err) {
return err
}
r.log.Error(err, "failed to get EndpointSlices", "namespace", string(*backendRef.Namespace),
backendRefKind, string(backendRef.Name))
logger.Error(err, "failed to get EndpointSlices")
} else {
for i := range endpointSliceList.Items {
endpointSlice := &endpointSliceList.Items[i]
key := utils.NamespacedName(endpointSlice).String()
if !resourceMappings.allAssociatedEndpointSlices.Has(key) {
resourceMappings.allAssociatedEndpointSlices.Insert(key)
r.log.Info("added EndpointSlice to resource tree",
logger.Info("added EndpointSlice to resource tree",
"namespace", endpointSlice.Namespace,
"name", endpointSlice.Name)
gwcResource.EndpointSlices = append(gwcResource.EndpointSlices, endpointSlice)
Expand Down
82 changes: 31 additions & 51 deletions internal/provider/kubernetes/kubernetes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1412,13 +1412,9 @@ func TestNamespaceSelectorProvider(t *testing.T) {
return resources.GatewayAPIResources.Len() != 0
}, defaultWait, defaultTick)

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
return res != nil && len(res.Gateways) == 1
}, defaultWait, defaultTick)
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.Gateways) == 1
})

_, ok := resources.GatewayStatuses.Load(types.NamespacedName{Name: "non-watched-gateway", Namespace: nonWatchedNS.Name})
require.Equal(t, false, ok)
Expand Down Expand Up @@ -1560,54 +1556,28 @@ func TestNamespaceSelectorProvider(t *testing.T) {
require.NoError(t, cli.Delete(ctx, nonWatchedUDPRoute))
}()

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
// The service number dependes on the service created and the backendRef
return res != nil && len(res.Services) == 5
}, defaultWait, defaultTick)

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
return res != nil && len(res.HTTPRoutes) == 1
}, defaultWait, defaultTick)
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.Services) == 1
})

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
return res != nil && len(res.TCPRoutes) == 1
}, defaultWait, defaultTick)
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.HTTPRoutes) == 1
})
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.TCPRoutes) == 1
})

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
return res != nil && len(res.TLSRoutes) == 1
}, defaultWait, defaultTick)
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.TLSRoutes) == 1
})

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
return res != nil && len(res.UDPRoutes) == 1
}, defaultWait, defaultTick)
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.UDPRoutes) == 1
})

require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gc.Name)
if !ok {
return false
}
return res != nil && len(res.GRPCRoutes) == 1
}, defaultWait, defaultTick)
requireResourceReady(t, resources, gc.Name, func(r *resource.Resources) bool {
return len(r.GRPCRoutes) == 1
})
}

func waitUntilGatewayClassResourcesAreReady(resources *message.ProviderResources, gatewayClassName string) (*resource.Resources, bool) {
Expand All @@ -1618,3 +1588,13 @@ func waitUntilGatewayClassResourcesAreReady(resources *message.ProviderResources

return res, true
}

func requireResourceReady(t *testing.T, resources *message.ProviderResources, gatewayClassName string, cmpFunc func(r *resource.Resources) bool) {
require.Eventually(t, func() bool {
res, ok := waitUntilGatewayClassResourcesAreReady(resources, gatewayClassName)
if !ok {
return false
}
return cmpFunc(res)
}, defaultWait, defaultTick)
}
3 changes: 3 additions & 0 deletions internal/provider/kubernetes/resource.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ type resourceMappings struct {
allAssociatedReferenceGrants sets.Set[string]
// Set for storing ServiceImports' NamespacedNames.
allAssociatedServiceImports sets.Set[string]
// Set for storing Services' NamespacedNames.
allAssociatedServices sets.Set[string]
// Set for storing EndpointSlices' NamespacedNames.
allAssociatedEndpointSlices sets.Set[string]
// Set for storing Backends' NamespacedNames.
Expand Down Expand Up @@ -75,6 +77,7 @@ func newResourceMapping() *resourceMappings {
allAssociatedGateways: sets.New[string](),
allAssociatedReferenceGrants: sets.New[string](),
allAssociatedServiceImports: sets.New[string](),
allAssociatedServices: sets.New[string](),
allAssociatedEndpointSlices: sets.New[string](),
allAssociatedBackends: sets.New[string](),
allAssociatedSecrets: sets.New[string](),
Expand Down
Loading