This is an automated email from the ASF dual-hosted git repository.

AlinsRan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/apisix-ingress-controller.git


The following commit(s) were added to refs/heads/master by this push:
     new a1c2ec17 fix: never treat an unreachable API server as a missing API 
resource (#2817)
a1c2ec17 is described below

commit a1c2ec172b75a10d5754b05c9b4700cd5360ca87
Author: AlinsRan <[email protected]>
AuthorDate: Tue Jul 28 14:48:30 2026 +0800

    fix: never treat an unreachable API server as a missing API resource (#2817)
---
 internal/controller/consumer_controller.go     |   9 +-
 internal/controller/gateway_controller.go      |  18 +-
 internal/controller/gatewayproxy_controller.go |   9 +-
 internal/controller/indexer/indexer.go         |  42 ++--
 internal/controller/tcproute_controller.go     |   6 +-
 internal/controller/tlsroute_controller.go     |   6 +-
 internal/controller/udproute_controller.go     |   6 +-
 internal/manager/controllers.go                | 100 +++++---
 internal/manager/run.go                        |  71 +++++-
 pkg/utils/k8s.go                               | 204 ++++++++++++++--
 pkg/utils/k8s_test.go                          | 324 +++++++++++++++++++++++++
 11 files changed, 705 insertions(+), 90 deletions(-)

diff --git a/internal/controller/consumer_controller.go 
b/internal/controller/consumer_controller.go
index 2c0a4eb0..0f265be5 100644
--- a/internal/controller/consumer_controller.go
+++ b/internal/controller/consumer_controller.go
@@ -62,7 +62,14 @@ type ConsumerReconciler struct { //nolint:revive
 
 // SetupWithManager sets up the controller with the Manager.
 func (r *ConsumerReconciler) SetupWithManager(mgr ctrl.Manager) error {
-       if config.ControllerConfig.DisableGatewayAPI || 
!pkgutils.HasAPIResource(mgr, &gatewayv1.Gateway{}) {
+       hasGatewayAPI := false
+       if !config.ControllerConfig.DisableGatewayAPI {
+               var err error
+               if hasGatewayAPI, err = pkgutils.HasAPIResource(mgr, 
&gatewayv1.Gateway{}); err != nil {
+                       return err
+               }
+       }
+       if !hasGatewayAPI {
                r.Log.Info("skipping Consumer controller setup as Gateway API 
is not available")
                return nil
        }
diff --git a/internal/controller/gateway_controller.go 
b/internal/controller/gateway_controller.go
index 5b5130d7..69cc89f5 100644
--- a/internal/controller/gateway_controller.go
+++ b/internal/controller/gateway_controller.go
@@ -105,19 +105,31 @@ func (r *GatewayReconciler) SetupWithManager(mgr 
ctrl.Manager) error {
                        
builder.WithPredicates(referenceGrantPredicates(KindGateway)),
                )
        }
-       if pkgutils.HasAPIResource(mgr, &gatewayv1alpha2.TCPRoute{}) {
+       hasTCPRoute, err := pkgutils.HasAPIResource(mgr, 
&gatewayv1alpha2.TCPRoute{})
+       if err != nil {
+               return err
+       }
+       if hasTCPRoute {
                bdr.Watches(
                        &gatewayv1alpha2.TCPRoute{},
                        
handler.EnqueueRequestsFromMapFunc(r.listGatewaysForStatusParentRefs),
                )
        }
-       if pkgutils.HasAPIResource(mgr, &gatewayv1alpha2.TLSRoute{}) {
+       hasTLSRoute, err := pkgutils.HasAPIResource(mgr, 
&gatewayv1alpha2.TLSRoute{})
+       if err != nil {
+               return err
+       }
+       if hasTLSRoute {
                bdr.Watches(
                        &gatewayv1alpha2.TLSRoute{},
                        
handler.EnqueueRequestsFromMapFunc(r.listGatewaysForStatusParentRefs),
                )
        }
-       if pkgutils.HasAPIResource(mgr, &gatewayv1alpha2.UDPRoute{}) {
+       hasUDPRoute, err := pkgutils.HasAPIResource(mgr, 
&gatewayv1alpha2.UDPRoute{})
+       if err != nil {
+               return err
+       }
+       if hasUDPRoute {
                bdr.Watches(
                        &gatewayv1alpha2.UDPRoute{},
                        
handler.EnqueueRequestsFromMapFunc(r.listGatewaysForStatusParentRefs),
diff --git a/internal/controller/gatewayproxy_controller.go 
b/internal/controller/gatewayproxy_controller.go
index 82bb2c71..ab53ad6f 100644
--- a/internal/controller/gatewayproxy_controller.go
+++ b/internal/controller/gatewayproxy_controller.go
@@ -54,9 +54,14 @@ type GatewayProxyController struct {
 }
 
 func (r *GatewayProxyController) SetupWithManager(mrg ctrl.Manager) error {
-       if config.ControllerConfig.DisableGatewayAPI || 
!pkgutils.HasAPIResource(mrg, &gatewayv1.Gateway{}) {
-               r.disableGatewayAPI = true
+       hasGatewayAPI := false
+       if !config.ControllerConfig.DisableGatewayAPI {
+               var err error
+               if hasGatewayAPI, err = pkgutils.HasAPIResource(mrg, 
&gatewayv1.Gateway{}); err != nil {
+                       return err
+               }
        }
+       r.disableGatewayAPI = !hasGatewayAPI
        builder := ctrl.NewControllerManagedBy(mrg).
                For(&v1alpha1.GatewayProxy{}).
                WithEventFilter(
diff --git a/internal/controller/indexer/indexer.go 
b/internal/controller/indexer/indexer.go
index 037e9e11..99f8b535 100644
--- a/internal/controller/indexer/indexer.go
+++ b/internal/controller/indexer/indexer.go
@@ -65,12 +65,16 @@ func SetupAPIv1alpha1Indexer(mgr ctrl.Manager) error {
                &v1alpha1.GatewayProxy{}:         setupGatewayProxyIndexer,
                &v1alpha1.L4RoutePolicy{}:        setupL4RoutePolicyIndexer,
        } {
-               if utils.HasAPIResource(mgr, resource) {
-                       if err := setup(mgr); err != nil {
-                               return err
-                       }
-               } else {
+               installed, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return err
+               }
+               if !installed {
                        setupLog.Info("Skipping indexer setup, API not found in 
cluster", "api", utils.FormatGVK(resource))
+                       continue
+               }
+               if err := setup(mgr); err != nil {
+                       return err
                }
        }
        return nil
@@ -87,12 +91,16 @@ func SetupAPIv2Indexer(mgr ctrl.Manager) error {
                &apiv2.ApisixTls{}:           setupApisixTlsIndexer,
                &apiv2.ApisixGlobalRule{}:    setupApisixGlobalRuleIndexer,
        } {
-               if utils.HasAPIResource(mgr, resource) {
-                       if err := setup(mgr); err != nil {
-                               return err
-                       }
-               } else {
+               installed, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return err
+               }
+               if !installed {
                        setupLog.Info("Skipping indexer setup, API not found in 
cluster", "api", utils.FormatGVK(resource))
+                       continue
+               }
+               if err := setup(mgr); err != nil {
+                       return err
                }
        }
        return nil
@@ -110,12 +118,16 @@ func SetupGatewayAPIIndexer(mgr ctrl.Manager) error {
                &gatewayv1alpha2.TLSRoute{}: setupTLSRouteIndexer,
                &gatewayv1.GatewayClass{}:   setupGatewayClassIndexer,
        } {
-               if utils.HasAPIResource(mgr, resource) {
-                       if err := setup(mgr); err != nil {
-                               return err
-                       }
-               } else {
+               installed, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return err
+               }
+               if !installed {
                        setupLog.Info("Skipping indexer setup, API not found in 
cluster", "api", utils.FormatGVK(resource))
+                       continue
+               }
+               if err := setup(mgr); err != nil {
+                       return err
                }
        }
        return nil
diff --git a/internal/controller/tcproute_controller.go 
b/internal/controller/tcproute_controller.go
index f76f6855..a2360157 100644
--- a/internal/controller/tcproute_controller.go
+++ b/internal/controller/tcproute_controller.go
@@ -101,7 +101,11 @@ func (r *TCPRouteReconciler) SetupWithManager(mgr 
ctrl.Manager) error {
 
        // L4RoutePolicy is an optional CRD. Only watch it when installed so the
        // controller still starts if the CRD has not been applied yet (e.g. 
upgrades).
-       r.supportsL4RoutePolicy = pkgutils.HasAPIResource(mgr, 
&v1alpha1.L4RoutePolicy{})
+       supportsL4RoutePolicy, err := pkgutils.HasAPIResource(mgr, 
&v1alpha1.L4RoutePolicy{})
+       if err != nil {
+               return err
+       }
+       r.supportsL4RoutePolicy = supportsL4RoutePolicy
        if r.supportsL4RoutePolicy {
                bdr.Watches(&v1alpha1.L4RoutePolicy{},
                        
handler.EnqueueRequestsFromMapFunc(r.listTCPRoutesForL4RoutePolicy),
diff --git a/internal/controller/tlsroute_controller.go 
b/internal/controller/tlsroute_controller.go
index dd6b7ff4..edb7381f 100644
--- a/internal/controller/tlsroute_controller.go
+++ b/internal/controller/tlsroute_controller.go
@@ -101,7 +101,11 @@ func (r *TLSRouteReconciler) SetupWithManager(mgr 
ctrl.Manager) error {
 
        // L4RoutePolicy is an optional CRD. Only watch it when installed so the
        // controller still starts if the CRD has not been applied yet (e.g. 
upgrades).
-       r.supportsL4RoutePolicy = pkgutils.HasAPIResource(mgr, 
&v1alpha1.L4RoutePolicy{})
+       supportsL4RoutePolicy, err := pkgutils.HasAPIResource(mgr, 
&v1alpha1.L4RoutePolicy{})
+       if err != nil {
+               return err
+       }
+       r.supportsL4RoutePolicy = supportsL4RoutePolicy
        if r.supportsL4RoutePolicy {
                bdr.Watches(&v1alpha1.L4RoutePolicy{},
                        
handler.EnqueueRequestsFromMapFunc(r.listTLSRoutesForL4RoutePolicy),
diff --git a/internal/controller/udproute_controller.go 
b/internal/controller/udproute_controller.go
index 33c698ae..17fb5f6a 100644
--- a/internal/controller/udproute_controller.go
+++ b/internal/controller/udproute_controller.go
@@ -101,7 +101,11 @@ func (r *UDPRouteReconciler) SetupWithManager(mgr 
ctrl.Manager) error {
 
        // L4RoutePolicy is an optional CRD. Only watch it when installed so the
        // controller still starts if the CRD has not been applied yet (e.g. 
upgrades).
-       r.supportsL4RoutePolicy = pkgutils.HasAPIResource(mgr, 
&v1alpha1.L4RoutePolicy{})
+       supportsL4RoutePolicy, err := pkgutils.HasAPIResource(mgr, 
&v1alpha1.L4RoutePolicy{})
+       if err != nil {
+               return err
+       }
+       r.supportsL4RoutePolicy = supportsL4RoutePolicy
        if r.supportsL4RoutePolicy {
                bdr.Watches(&v1alpha1.L4RoutePolicy{},
                        
handler.EnqueueRequestsFromMapFunc(r.listUDPRoutesForL4RoutePolicy),
diff --git a/internal/manager/controllers.go b/internal/manager/controllers.go
index 5e48b397..af518e3c 100644
--- a/internal/manager/controllers.go
+++ b/internal/manager/controllers.go
@@ -211,11 +211,15 @@ func setupGatewayAPIControllers(ctx context.Context, mgr 
manager.Manager, pro pr
                        Readier:  readier,
                },
        } {
-               if utils.HasAPIResource(mgr, resource) {
-                       runnables = append(runnables, controller)
-               } else {
-                       setupLog.Info("Skipping indexer setup, API not found in 
cluster", "api", utils.FormatGVK(resource))
+               installed, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return nil, err
                }
+               if !installed {
+                       setupLog.Info("Skipping controller setup, API not found 
in cluster", "api", utils.FormatGVK(resource))
+                       continue
+               }
+               runnables = append(runnables, controller)
        }
        return runnables, nil
 }
@@ -288,69 +292,88 @@ func setupAPIv2Controllers(ctx context.Context, mgr 
manager.Manager, pro provide
                        Updater: updater,
                },
        } {
-               if utils.HasAPIResource(mgr, resource) {
-                       runnables = append(runnables, controller)
-               } else {
-                       setupLog.Info("Skipping indexer setup, API not found in 
cluster", "api", utils.FormatGVK(resource))
+               installed, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return nil, err
+               }
+               if !installed {
+                       setupLog.Info("Skipping controller setup, API not found 
in cluster", "api", utils.FormatGVK(resource))
+                       continue
                }
+               runnables = append(runnables, controller)
        }
        return runnables, nil
 }
 
-func registerReadiness(mgr manager.Manager, readier 
readiness.ReadinessManager) {
+func registerReadiness(mgr manager.Manager, readier 
readiness.ReadinessManager) error {
        log := ctrl.LoggerFrom(context.Background()).WithName("readiness")
 
-       registerAPIv2ForReadiness(mgr, log, readier)
+       if err := registerAPIv2ForReadiness(mgr, log, readier); err != nil {
+               return err
+       }
        if !config.ControllerConfig.DisableGatewayAPI {
-               registerGatewayAPIForReadiness(mgr, log, readier)
+               if err := registerGatewayAPIForReadiness(mgr, log, readier); 
err != nil {
+                       return err
+               }
        }
-       registerAPIv1alpha1ForReadiness(mgr, log, readier)
+       return registerAPIv1alpha1ForReadiness(mgr, log, readier)
 }
 
 func registerGatewayAPIForReadiness(
        mgr manager.Manager,
        log logr.Logger,
        readier readiness.ReadinessManager,
-) {
-       var installed []schema.GroupVersionKind
-       for _, resource := range []client.Object{
+) error {
+       resources := []client.Object{
                &gatewayv1.HTTPRoute{},
                &gatewayv1.GRPCRoute{},
                &gatewayv1alpha2.TCPRoute{},
                &gatewayv1alpha2.UDPRoute{},
                &gatewayv1alpha2.TLSRoute{},
-       } {
+       }
+       installed := make([]schema.GroupVersionKind, 0, len(resources))
+       for _, resource := range resources {
                gvk := types.GvkOf(resource)
-               if utils.HasAPIResource(mgr, resource) {
-                       installed = append(installed, gvk)
-               } else {
+               has, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return err
+               }
+               if !has {
                        log.Info("Skipping readiness registration, API not 
found", "gvk", gvk)
+                       continue
                }
+               installed = append(installed, gvk)
        }
        if len(installed) == 0 {
-               return
+               return nil
        }
 
        readier.RegisterGVK(readiness.GVKConfig{GVKs: installed})
+       return nil
 }
 
 func registerAPIv2ForReadiness(
        mgr manager.Manager,
        log logr.Logger,
        readier readiness.ReadinessManager,
-) {
-       var installed []schema.GroupVersionKind
-       for _, resource := range apiV2ReadinessResources() {
+) error {
+       resources := apiV2ReadinessResources()
+       installed := make([]schema.GroupVersionKind, 0, len(resources))
+       for _, resource := range resources {
                gvk := types.GvkOf(resource)
-               if utils.HasAPIResource(mgr, resource) {
-                       installed = append(installed, gvk)
-               } else {
+               has, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return err
+               }
+               if !has {
                        log.Info("Skipping readiness registration, API not 
found", "gvk", gvk)
+                       continue
                }
+               installed = append(installed, gvk)
        }
 
        if len(installed) == 0 {
-               return
+               return nil
        }
 
        readier.RegisterGVK(readiness.GVKConfig{
@@ -361,6 +384,7 @@ func registerAPIv2ForReadiness(
                        return ingressClass != nil
                }),
        })
+       return nil
 }
 
 func apiV2ReadinessResources() []client.Object {
@@ -377,20 +401,25 @@ func registerAPIv1alpha1ForReadiness(
        mgr manager.Manager,
        log logr.Logger,
        readier readiness.ReadinessManager,
-) {
-       var installed []schema.GroupVersionKind
-       for _, resource := range []client.Object{
+) error {
+       resources := []client.Object{
                &v1alpha1.Consumer{},
-       } {
+       }
+       installed := make([]schema.GroupVersionKind, 0, len(resources))
+       for _, resource := range resources {
                gvk := types.GvkOf(resource)
-               if utils.HasAPIResource(mgr, resource) {
-                       installed = append(installed, gvk)
-               } else {
+               has, err := utils.HasAPIResource(mgr, resource)
+               if err != nil {
+                       return err
+               }
+               if !has {
                        log.Info("Skipping readiness registration, API not 
found", "gvk", gvk)
+                       continue
                }
+               installed = append(installed, gvk)
        }
        if len(installed) == 0 {
-               return
+               return nil
        }
 
        readier.RegisterGVK(readiness.GVKConfig{
@@ -403,4 +432,5 @@ func registerAPIv1alpha1ForReadiness(
                        return 
controller.MatchConsumerGatewayRef(context.Background(), mgr.GetClient(), log, 
consumer)
                }),
        })
+       return nil
 }
diff --git a/internal/manager/run.go b/internal/manager/run.go
index 315644da..34c2d338 100644
--- a/internal/manager/run.go
+++ b/internal/manager/run.go
@@ -24,7 +24,6 @@ import (
 
        "github.com/go-logr/logr"
        "k8s.io/apimachinery/pkg/runtime"
-       "k8s.io/apimachinery/pkg/runtime/schema"
        utilruntime "k8s.io/apimachinery/pkg/util/runtime"
        "k8s.io/apimachinery/pkg/util/version"
        "k8s.io/client-go/discovery"
@@ -49,6 +48,7 @@ import (
        "github.com/apache/apisix-ingress-controller/internal/provider"
        _ "github.com/apache/apisix-ingress-controller/internal/provider/init"
        _ "github.com/apache/apisix-ingress-controller/pkg/metrics"
+       "github.com/apache/apisix-ingress-controller/pkg/utils"
 )
 
 var (
@@ -86,6 +86,11 @@ func Run(ctx context.Context, logger logr.Logger) error {
 
        setupLog := ctrl.LoggerFrom(ctx).WithName("setup")
 
+       // SetupSignalHandler must be called exactly once. Install it before 
setup
+       // starts so that a shutdown signal is honored while waiting for the API
+       // server, not only once the manager is running.
+       signalCtx := ctrl.SetupSignalHandler()
+
        // if the enable-http2 flag is false (the default), http/2 should be 
disabled
        // due to its vulnerabilities. More specifically, disabling http/2 will
        // prevent from being vulnerable to the HTTP/2 Stream Cancellation and
@@ -171,11 +176,40 @@ func Run(ctx context.Context, logger logr.Logger) error {
                return err
        }
 
+       // Which API resources are installed is detected once, below, and gates 
the
+       // registration of field indexes, controllers and readiness checks for 
the
+       // whole lifetime of the process. Detecting them against an unreachable 
API
+       // server would silently classify everything as "not installed" and 
leave the
+       // controller permanently degraded until it is restarted, so wait for 
the API
+       // server to answer first.
+       if err := utils.WaitForAPIServer(signalCtx, mgr.GetConfig(), setupLog); 
err != nil {
+               setupLog.Error(err, "unable to reach the Kubernetes API server")
+               return err
+       }
+
+       // API resource detection runs outside the manager and does not observe
+       // signalCtx, so check between setup phases: now that the signal 
handler is
+       // installed this early, a shutdown signal would otherwise be swallowed 
until
+       // mgr.Start.
+       checkShutdown := func() error {
+               if err := signalCtx.Err(); err != nil {
+                       setupLog.Info("shutdown requested during setup, 
stopping", "reason", err)
+                       return err
+               }
+               return nil
+       }
+
        // Check Kubernetes cluster version
        checkK8sVersion(mgr, setupLog)
 
        readier := readiness.NewReadinessManager(mgr.GetClient(), logger)
-       registerReadiness(mgr, readier)
+       if err := registerReadiness(mgr, readier); err != nil {
+               setupLog.Error(err, "unable to register readiness checks")
+               return err
+       }
+       if err := checkShutdown(); err != nil {
+               return err
+       }
 
        if err := mgr.Add(readier); err != nil {
                setupLog.Error(err, "unable to add readiness manager")
@@ -214,16 +248,26 @@ func Run(ctx context.Context, logger logr.Logger) error {
                return err
        }
 
-       setupLog.Info("check ReferenceGrants is enabled")
-       _, err = mgr.GetRESTMapper().KindsFor(schema.GroupVersionResource{
-               Group:    v1beta1.GroupVersion.Group,
-               Version:  v1beta1.GroupVersion.Version,
-               Resource: "referencegrants",
-       })
-       if err != nil {
-               setupLog.Info("CRD ReferenceGrants is not installed", "err", 
err)
+       // ReferenceGrant is a Gateway API kind and only consulted by Gateway 
API
+       // paths, so skip the detection entirely when Gateway API is disabled.
+       hasReferenceGrant := false
+       if config.ControllerConfig.DisableGatewayAPI {
+               setupLog.Info("Gateway API is disabled, skipping the 
ReferenceGrants check")
+       } else {
+               setupLog.Info("check ReferenceGrants is enabled")
+               if hasReferenceGrant, err = utils.HasAPIResource(mgr, 
&v1beta1.ReferenceGrant{}); err != nil {
+                       setupLog.Error(err, "unable to detect whether 
ReferenceGrants is installed")
+                       return err
+               }
+               if !hasReferenceGrant {
+                       setupLog.Info("CRD ReferenceGrants is not installed, 
cross-namespace references will be rejected",
+                               "gvk", 
utils.FormatGVK(&v1beta1.ReferenceGrant{}))
+               }
+       }
+       controller.SetEnableReferenceGrant(hasReferenceGrant)
+       if err := checkShutdown(); err != nil {
+               return err
        }
-       controller.SetEnableReferenceGrant(err == nil)
 
        setupLog.Info("setting up controllers")
        controllers, err := setupControllers(ctx, mgr, provider, 
updater.Writer(), readier)
@@ -236,6 +280,9 @@ func Run(ctx context.Context, logger logr.Logger) error {
                if err := c.SetupWithManager(mgr); err != nil {
                        return err
                }
+               if err := checkShutdown(); err != nil {
+                       return err
+               }
        }
 
        // +kubebuilder:scaffold:builder
@@ -263,7 +310,7 @@ func Run(ctx context.Context, logger logr.Logger) error {
        }
 
        setupLog.Info("starting controller manager")
-       return mgr.Start(ctrl.SetupSignalHandler())
+       return mgr.Start(signalCtx)
 }
 
 func checkK8sVersion(mgr ctrl.Manager, logger logr.Logger) {
diff --git a/pkg/utils/k8s.go b/pkg/utils/k8s.go
index 425f2ff9..97c46e2b 100644
--- a/pkg/utils/k8s.go
+++ b/pkg/utils/k8s.go
@@ -18,8 +18,18 @@
 package utils
 
 import (
+       "context"
+       "errors"
+       "fmt"
+       "time"
+
        "github.com/go-logr/logr"
+       apierrors "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime/schema"
+       "k8s.io/apimachinery/pkg/util/wait"
        "k8s.io/client-go/discovery"
+       "k8s.io/client-go/rest"
        ctrl "sigs.k8s.io/controller-runtime"
        "sigs.k8s.io/controller-runtime/pkg/client"
        "sigs.k8s.io/controller-runtime/pkg/client/apiutil"
@@ -27,22 +37,81 @@ import (
        "github.com/apache/apisix-ingress-controller/internal/types"
 )
 
-// HasAPIResource checks if a specific API resource is available in the 
current cluster.
-// It uses the Discovery API to query the cluster's available resources and 
returns true
-// if the resource is found, false otherwise.
-func HasAPIResource(mgr ctrl.Manager, obj client.Object) bool {
+const (
+       // A discovery request that gives no definitive answer is retried a few 
times
+       // to absorb short blips. A longer outage is reported to the caller 
instead:
+       // it is not this function's job to decide how long to wait for the API 
server.
+       discoveryRetryInitialInterval = 500 * time.Millisecond
+       discoveryRetrySteps           = 5 // ~7.5s, up to ~8.3s with jitter
+
+       // The startup wait for the API server is deliberately finite. The 
health
+       // probe endpoint is only served once mgr.Start runs, so a pod still 
waiting
+       // here fails its liveness probe and is killed anyway; giving up with an
+       // explicit error and letting the pod restart says more in the logs than
+       // being killed mid-wait. Longer outages are absorbed by the restart 
backoff.
+       apiServerWaitInitialInterval = 1 * time.Second
+       apiServerWaitSteps           = 6 // ~31s, up to ~34s with jitter
+
+       backoffFactor = 2
+       backoffJitter = 0.1
+)
+
+// Neither backoff sets Cap: wait.Backoff zeroes the remaining steps as soon as
+// the capped interval is reached, which would silently cut the retries short.
+
+func discoveryBackoff() wait.Backoff {
+       return wait.Backoff{
+               Duration: discoveryRetryInitialInterval,
+               Factor:   backoffFactor,
+               Jitter:   backoffJitter,
+               Steps:    discoveryRetrySteps,
+       }
+}
+
+func apiServerWaitBackoff() wait.Backoff {
+       return wait.Backoff{
+               Duration: apiServerWaitInitialInterval,
+               Factor:   backoffFactor,
+               Jitter:   backoffJitter,
+               Steps:    apiServerWaitSteps,
+       }
+}
+
+// HasAPIResource reports whether the API resource of obj is served by the 
cluster.
+//
+// Only a 404 from discovery reports the resource as absent. Every other 
failure
+// returns an error, including a Forbidden discovery request: callers must not
+// fall back to "resource absent" there. Detection runs once at startup and 
gates
+// field index, controller and readiness registration, so an API server outage 
or
+// an RBAC mistake would otherwise leave the controller permanently degraded
+// until it is restarted.
+func HasAPIResource(mgr ctrl.Manager, obj client.Object) (bool, error) {
        return HasAPIResourceWithLogger(mgr, obj, 
ctrl.Log.WithName("api-detection"))
 }
 
 // HasAPIResourceWithLogger is the same as HasAPIResource but accepts a custom 
logger
 // for more detailed debugging information.
-func HasAPIResourceWithLogger(mgr ctrl.Manager, obj client.Object, logger 
logr.Logger) bool {
+func HasAPIResourceWithLogger(mgr ctrl.Manager, obj client.Object, logger 
logr.Logger) (bool, error) {
        gvk, err := apiutil.GVKForObject(obj, mgr.GetScheme())
        if err != nil {
-               logger.Info("cannot derive GVK from scheme", "error", err)
-               return false
+               return false, fmt.Errorf("cannot derive GVK from scheme: %w", 
err)
+       }
+
+       // Create discovery client
+       discoveryClient, err := 
discovery.NewDiscoveryClientForConfig(mgr.GetConfig())
+       if err != nil {
+               return false, fmt.Errorf("failed to create discovery client: 
%w", err)
        }
 
+       return hasAPIResource(discoveryClient, gvk, discoveryBackoff(), logger)
+}
+
+func hasAPIResource(
+       discoveryClient discovery.DiscoveryInterface,
+       gvk schema.GroupVersionKind,
+       backoff wait.Backoff,
+       logger logr.Logger,
+) (bool, error) {
        groupVersion := gvk.GroupVersion().String()
 
        logger = logger.WithValues(
@@ -52,29 +121,126 @@ func HasAPIResourceWithLogger(mgr ctrl.Manager, obj 
client.Object, logger logr.L
                "groupVersion", groupVersion,
        )
 
-       // Create discovery client
-       discoveryClient, err := 
discovery.NewDiscoveryClientForConfig(mgr.GetConfig())
-       if err != nil {
-               logger.Info("failed to create discovery client", "error", err)
-               return false
-       }
-
        // Query server resources for the specific group/version
-       apiResources, err := 
discoveryClient.ServerResourcesForGroupVersion(groupVersion)
-       if err != nil {
+       var apiResources *metav1.APIResourceList
+       err := retryUntilDefinitive(backoff, logger, func() error {
+               var err error
+               apiResources, err = 
discoveryClient.ServerResourcesForGroupVersion(groupVersion)
+               return err
+       })
+       switch {
+       case err == nil:
+       case apierrors.IsNotFound(err):
                logger.Info("group/version not available in cluster", "error", 
err)
-               return false
+               return false, nil
+       case apierrors.IsForbidden(err):
+               // Definitive, but no evidence of absence: the resource may 
well be served
+               // and only discovery is denied. Resolving it to "absent" would 
be cached
+               // for the lifetime of the manager and permanently skip the 
controller, so
+               // an RBAC misconfiguration must abort startup instead.
+               return false, fmt.Errorf("discovery of %s is forbidden, check 
the discovery "+
+                       "permissions of the controller service account: %w", 
gvk, err)
+       default:
+               return false, fmt.Errorf("failed to detect API resource %s: 
%w", gvk, err)
        }
 
        // Check if the specific kind exists in the resource list
        for _, res := range apiResources.APIResources {
                if res.Kind == gvk.Kind {
-                       return true
+                       return true, nil
                }
        }
 
        logger.Info("API resource kind not found in group/version")
-       return false
+       return false, nil
+}
+
+// WaitForAPIServer blocks until the API server answers a discovery request, 
ctx
+// is done, or the wait budget is exhausted.
+//
+// Capability detection (field indexes, optional CRDs, ReferenceGrant support) 
is
+// decided once at startup and cannot be revised afterwards, so it must not run
+// against an unreachable API server.
+func WaitForAPIServer(ctx context.Context, cfg *rest.Config, logger 
logr.Logger) error {
+       discoveryClient, err := discovery.NewDiscoveryClientForConfig(cfg)
+       if err != nil {
+               return fmt.Errorf("failed to create discovery client: %w", err)
+       }
+       return waitForAPIServer(ctx, func() error {
+               _, err := discoveryClient.ServerVersion()
+               return err
+       }, apiServerWaitBackoff(), logger)
+}
+
+func waitForAPIServer(ctx context.Context, probe func() error, backoff 
wait.Backoff, logger logr.Logger) error {
+       var lastErr error
+       if err := wait.ExponentialBackoffWithContext(ctx, backoff, 
func(context.Context) (bool, error) {
+               if lastErr = probe(); !apiServerAnswered(lastErr) {
+                       logger.Info("waiting for the Kubernetes API server to 
become reachable", "error", lastErr)
+                       return false, nil
+               }
+               if lastErr != nil {
+                       logger.Info("the Kubernetes API server rejected the 
probe but is reachable, continuing", "error", lastErr)
+               }
+               return true, nil
+       }); err != nil {
+               // ExponentialBackoffWithContext only ever returns ctx.Err() or 
its own
+               // interrupted error, so this tells cancellation and an 
exhausted budget
+               // apart without racing against ctx.
+               if !errors.Is(err, context.Canceled) && !errors.Is(err, 
context.DeadlineExceeded) && lastErr != nil {
+                       err = lastErr
+               }
+               return fmt.Errorf("give up waiting for the Kubernetes API 
server: %w", err)
+       }
+       return nil
+}
+
+// apiServerAnswered reports whether err still proves the API server responded.
+// An authn/authz rejection is an answer: a hardened cluster may not grant the
+// controller's service account access to the probed endpoint, and that is no
+// reason to hold up startup.
+func apiServerAnswered(err error) bool {
+       return err == nil || apierrors.IsUnauthorized(err) || 
apierrors.IsForbidden(err)
+}
+
+// retryUntilDefinitive runs probe until it returns nil or a definitive answer,
+// and reports the last failure once the retry budget is exhausted.
+func retryUntilDefinitive(backoff wait.Backoff, logger logr.Logger, probe 
func() error) error {
+       var probed bool
+       var lastErr error
+       err := wait.ExponentialBackoff(backoff, func() (bool, error) {
+               probed = true
+               if lastErr = probe(); lastErr == nil || isDefinitive(lastErr) {
+                       return true, nil
+               }
+               logger.Info("discovery request failed, retrying", "error", 
lastErr)
+               return false, nil
+       })
+       if !probed {
+               // A backoff with no steps left never runs the condition. 
Reporting nil
+               // here would be read as "the resource is present" by every 
caller.
+               return fmt.Errorf("discovery was never attempted: %w", err)
+       }
+       return lastErr
+}
+
+// isDefinitive reports whether err settles the discovery request, so that
+// retrying it cannot change the outcome. Only the caller decides what a
+// definitive failure means: NotFound is absence, Forbidden is a 
misconfiguration.
+func isDefinitive(err error) bool {
+       switch {
+       // Discovery of a group/version that is not installed answers 404.
+       case apierrors.IsNotFound(err):
+               return true
+       // RBAC forbids discovery: waiting will not change that, but it says 
nothing
+       // about whether the resource is served.
+       case apierrors.IsForbidden(err):
+               return true
+       // Connection refused, timeouts, EOF, 5xx, throttling: the API server 
may
+       // answer differently once it is reachable again, so no conclusion yet.
+       default:
+               return false
+       }
 }
 
 func FormatGVK(obj client.Object) string {
diff --git a/pkg/utils/k8s_test.go b/pkg/utils/k8s_test.go
new file mode 100644
index 00000000..73ce3579
--- /dev/null
+++ b/pkg/utils/k8s_test.go
@@ -0,0 +1,324 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package utils
+
+import (
+       "context"
+       "errors"
+       "net"
+       "net/http"
+       "net/http/httptest"
+       "testing"
+       "time"
+
+       "github.com/go-logr/logr"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+       apierrors "k8s.io/apimachinery/pkg/api/errors"
+       "k8s.io/apimachinery/pkg/runtime/schema"
+       "k8s.io/apimachinery/pkg/util/wait"
+       "k8s.io/client-go/discovery"
+       "k8s.io/client-go/rest"
+)
+
+// fastBackoff keeps the tests quick while preserving the shape of the
+// production backoffs (a growing interval over a fixed number of steps).
+func fastBackoff(steps int) wait.Backoff {
+       return wait.Backoff{Duration: time.Millisecond, Factor: backoffFactor, 
Jitter: backoffJitter, Steps: steps}
+}
+
+// declaredBudget is how long a backoff sleeps in total before jitter, i.e. the
+// sum of the intervals between its attempts.
+func declaredBudget(backoff wait.Backoff) time.Duration {
+       var total, interval time.Duration = 0, backoff.Duration
+       for i := 0; i < backoff.Steps-1; i++ {
+               total += interval
+               interval = time.Duration(float64(interval) * backoff.Factor)
+       }
+       return total
+}
+
+// countingBackoff reports how many times a backoff actually runs its 
condition.
+func countAttempts(backoff wait.Backoff) int {
+       var attempts int
+       _ = wait.ExponentialBackoff(backoff, func() (bool, error) {
+               attempts++
+               return false, nil
+       })
+       return attempts
+}
+
+// TestBackoffsUseAllTheirSteps guards against a subtle wait.Backoff behavior:
+// setting Cap zeroes the remaining steps as soon as the capped interval is
+// reached, silently cutting the retries short. Both backoffs must run every
+// step they declare.
+func TestBackoffsUseAllTheirSteps(t *testing.T) {
+       for _, tt := range []struct {
+               name    string
+               backoff wait.Backoff
+               steps   int
+               budget  time.Duration // as documented next to the constants
+       }{
+               {"discovery", discoveryBackoff(), discoveryRetrySteps, 7500 * 
time.Millisecond},
+               {"apiServerWait", apiServerWaitBackoff(), apiServerWaitSteps, 
31 * time.Second},
+       } {
+               t.Run(tt.name, func(t *testing.T) {
+                       assert.Zero(t, tt.backoff.Cap, "Cap truncates the 
retries, see wait.delay()")
+                       assert.Equal(t, tt.budget, declaredBudget(tt.backoff), 
"the documented budget must match the backoff")
+
+                       b := tt.backoff
+                       b.Duration = time.Microsecond // keep the test fast, 
preserve the shape
+                       assert.Equal(t, tt.steps, countAttempts(b))
+               })
+       }
+}
+
+func TestIsDefinitive(t *testing.T) {
+       gr := schema.GroupResource{Group: "gateway.networking.k8s.io", 
Resource: "tcproutes"}
+       tests := []struct {
+               name       string
+               err        error
+               definitive bool
+       }{
+               // Definitive: the group/version is genuinely not served.
+               {"not found", apierrors.NewNotFound(gr, "tcproutes"), true},
+               // Definitive: RBAC will not self-heal by waiting.
+               {"forbidden", apierrors.NewForbidden(gr, "tcproutes", 
errors.New("nope")), true},
+               // Indeterminate: the API server is temporarily unreachable or 
overloaded.
+               {"connection refused", &net.OpError{Op: "dial", Err: 
errors.New("connection refused")}, false},
+               {"service unavailable", 
apierrors.NewServiceUnavailable("apiserver is starting up"), false},
+               {"server timeout", apierrors.NewServerTimeout(gr, "get", 1), 
false},
+               {"too many requests", apierrors.NewTooManyRequestsError("slow 
down"), false},
+               {"internal error", 
apierrors.NewInternalError(errors.New("boom")), false},
+               {"generic transport error", errors.New("EOF"), false},
+       }
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       assert.Equal(t, tt.definitive, isDefinitive(tt.err))
+               })
+       }
+}
+
+// TestRetryUntilDefinitive_ReportsIndeterminateFailure is the core regression
+// test for #2734: a discovery failure that gives no definitive answer must
+// surface as an error so the caller aborts startup, instead of collapsing into
+// "resource absent" and permanently skipping the field index and its 
controller.
+func TestRetryUntilDefinitive_ReportsIndeterminateFailure(t *testing.T) {
+       var calls int
+       probe := func() error {
+               calls++
+               return apierrors.NewServiceUnavailable("apiserver not ready")
+       }
+
+       err := retryUntilDefinitive(fastBackoff(3), logr.Discard(), probe)
+       require.Error(t, err, "an unreachable API server must not be reported 
as a definitive answer")
+       assert.False(t, isDefinitive(err))
+       assert.Equal(t, 3, calls, "the whole retry budget must be used")
+}
+
+func TestRetryUntilDefinitive_TransientThenSuccess(t *testing.T) {
+       var calls int
+       probe := func() error {
+               calls++
+               if calls < 3 {
+                       return apierrors.NewServiceUnavailable("apiserver not 
ready")
+               }
+               return nil
+       }
+
+       require.NoError(t, retryUntilDefinitive(fastBackoff(5), logr.Discard(), 
probe))
+       assert.Equal(t, 3, calls, "an indeterminate failure must be retried, 
not given up on")
+}
+
+// TestRetryUntilDefinitive_DefinitiveAbsent ensures a genuinely absent CRD is
+// answered immediately without retrying, preserving optional-CRD support.
+func TestRetryUntilDefinitive_DefinitiveAbsent(t *testing.T) {
+       gr := schema.GroupResource{Group: "gateway.networking.k8s.io", 
Resource: "tcproutes"}
+       var calls int
+       probe := func() error {
+               calls++
+               return apierrors.NewNotFound(gr, "tcproutes")
+       }
+
+       err := retryUntilDefinitive(fastBackoff(5), logr.Discard(), probe)
+       assert.True(t, isDefinitive(err))
+       assert.Equal(t, 1, calls, "a definitive NotFound must not be retried")
+}
+
+// TestRetryUntilDefinitive_NeverProbed guards the caller's assumption that a
+// nil error means "the probe succeeded": a backoff with no steps left never 
runs
+// the condition, and reporting nil there would be read as "resource present".
+func TestRetryUntilDefinitive_NeverProbed(t *testing.T) {
+       var calls int
+       probe := func() error {
+               calls++
+               return nil
+       }
+
+       err := retryUntilDefinitive(wait.Backoff{Steps: 0}, logr.Discard(), 
probe)
+       require.Error(t, err)
+       assert.Zero(t, calls)
+       assert.False(t, isDefinitive(err), "an unattempted probe must not look 
like a definitive answer either")
+}
+
+func TestWaitForAPIServer_RetriesUntilReachable(t *testing.T) {
+       var calls int
+       probe := func() error {
+               calls++
+               if calls < 3 {
+                       return &net.OpError{Op: "dial", Err: 
errors.New("connection refused")}
+               }
+               return nil
+       }
+
+       require.NoError(t, waitForAPIServer(context.Background(), probe, 
fastBackoff(6), logr.Discard()))
+       assert.Equal(t, 3, calls)
+}
+
+// TestWaitForAPIServer_AuthRejectionIsReachable covers a hardened cluster that
+// does not grant the controller access to the probed endpoint: the API server
+// answered, so startup must proceed instead of stalling until the budget runs
+// out and then failing on a perfectly healthy cluster.
+func TestWaitForAPIServer_AuthRejectionIsReachable(t *testing.T) {
+       for name, probeErr := range map[string]error{
+               "unauthorized": apierrors.NewUnauthorized("no token"),
+               "forbidden":    apierrors.NewForbidden(schema.GroupResource{}, 
"", errors.New("nope")),
+       } {
+               t.Run(name, func(t *testing.T) {
+                       var calls int
+                       probe := func() error {
+                               calls++
+                               return probeErr
+                       }
+                       require.NoError(t, 
waitForAPIServer(context.Background(), probe, fastBackoff(5), logr.Discard()))
+                       assert.Equal(t, 1, calls, "an answer, even a rejection, 
must not be retried")
+               })
+       }
+}
+
+func TestWaitForAPIServer_GivesUpWithTheLastError(t *testing.T) {
+       probe := func() error { return 
apierrors.NewServiceUnavailable("apiserver not ready") }
+
+       err := waitForAPIServer(context.Background(), probe, fastBackoff(3), 
logr.Discard())
+       require.Error(t, err)
+       assert.True(t, apierrors.IsServiceUnavailable(errors.Unwrap(err)), "the 
last probe error must be reported, got %v", err)
+}
+
+// TestWaitForAPIServer_HonorsContext ensures the wait is interrupted 
mid-flight,
+// so a shutdown signal is not ignored while the API server is unreachable.
+func TestWaitForAPIServer_HonorsContext(t *testing.T) {
+       ctx, cancel := context.WithCancel(context.Background())
+       defer cancel()
+
+       var calls int
+       probe := func() error {
+               if calls++; calls == 2 {
+                       cancel()
+               }
+               return &net.OpError{Op: "dial", Err: errors.New("connection 
refused")}
+       }
+
+       // A budget far larger than the number of attempts the test expects: 
the wait
+       // must end because ctx was cancelled, not because it ran out of steps.
+       err := waitForAPIServer(ctx, probe, fastBackoff(100), logr.Discard())
+       require.ErrorIs(t, err, context.Canceled)
+       assert.Equal(t, 2, calls, "the probe must stop being called once ctx is 
cancelled")
+}
+
+// newDiscoveryClient serves discovery responses from handler.
+func newDiscoveryClient(t *testing.T, handler http.HandlerFunc) 
discovery.DiscoveryInterface {
+       t.Helper()
+       srv := httptest.NewServer(handler)
+       t.Cleanup(srv.Close)
+       c, err := discovery.NewDiscoveryClientForConfig(&rest.Config{Host: 
srv.URL})
+       require.NoError(t, err)
+       return c
+}
+
+// TestHasAPIResource covers the contract every caller depends on, against a 
real
+// discovery client rather than a fake: only a definitive answer may resolve to
+// "absent", everything else must be an error.
+func TestHasAPIResource(t *testing.T) {
+       gvk := schema.GroupVersionKind{Group: "gateway.networking.k8s.io", 
Version: "v1alpha2", Kind: "TCPRoute"}
+
+       tests := []struct {
+               name      string
+               handler   http.HandlerFunc
+               wantFound bool
+               wantErr   bool
+       }{
+               {
+                       name: "kind served",
+                       handler: func(w http.ResponseWriter, _ *http.Request) {
+                               w.Header().Set("Content-Type", 
"application/json")
+                               _, _ = 
w.Write([]byte(`{"kind":"APIResourceList","groupVersion":"gateway.networking.k8s.io/v1alpha2",`
 +
+                                       
`"resources":[{"name":"tcproutes","kind":"TCPRoute"}]}`))
+                       },
+                       wantFound: true,
+               },
+               {
+                       name: "group/version served but kind missing",
+                       handler: func(w http.ResponseWriter, _ *http.Request) {
+                               w.Header().Set("Content-Type", 
"application/json")
+                               _, _ = 
w.Write([]byte(`{"kind":"APIResourceList","groupVersion":"gateway.networking.k8s.io/v1alpha2",`
 +
+                                       
`"resources":[{"name":"udproutes","kind":"UDPRoute"}]}`))
+                       },
+               },
+               {
+                       // What a real API server returns for a group/version 
that is not
+                       // installed: a plain-text 404, not a Status object.
+                       name: "group/version not installed",
+                       handler: func(w http.ResponseWriter, _ *http.Request) {
+                               http.NotFound(w, nil)
+                       },
+               },
+               {
+                       // The #2734 scenario: the API server cannot answer. 
This must not be
+                       // reported as "the CRD is not installed".
+                       name: "api server unavailable",
+                       handler: func(w http.ResponseWriter, _ *http.Request) {
+                               w.WriteHeader(http.StatusServiceUnavailable)
+                       },
+                       wantErr: true,
+               },
+               {
+                       // Forbidden is definitive but proves nothing about the 
resource: the
+                       // CRD may well be served and only discovery denied. 
Resolving it to
+                       // "absent" would permanently skip the controller, so 
it must abort
+                       // startup and point at the RBAC instead.
+                       name: "discovery forbidden",
+                       handler: func(w http.ResponseWriter, _ *http.Request) {
+                               w.WriteHeader(http.StatusForbidden)
+                       },
+                       wantErr: true,
+               },
+       }
+
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       found, err := hasAPIResource(newDiscoveryClient(t, 
tt.handler), gvk, fastBackoff(2), logr.Discard())
+                       if tt.wantErr {
+                               require.Error(t, err)
+                               assert.False(t, found)
+                               return
+                       }
+                       require.NoError(t, err)
+                       assert.Equal(t, tt.wantFound, found)
+               })
+       }
+}


Reply via email to